Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -1439,8 +1439,6 @@ class BeamModulePlugin implements Plugin<Project> {
include 'src/*/java/**/*.java'
exclude '**/DefaultPackageTest.java'
}
// For spotless:off and spotless:on
toggleOffOn()
}
}

Expand Down

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,6 @@
*/
package org.apache.beam.sdk.schemas;

import static org.apache.beam.sdk.util.Preconditions.checkStateNotNull;

import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
Expand Down Expand Up @@ -407,17 +405,12 @@ Object convert(OneOfType.Value value) {

@NonNull
FieldValueGetter<@NonNull Object, Object> converter =
checkStateNotNull(
Verify.verifyNotNull(
converters.get(caseType.getValue()),
"Missing OneOf converter for case %s.",
caseType);

Object convertedValue =
checkStateNotNull(
converter.get(value.getValue()),
"Bug! converting a non-null value in a OneOf resulted in null result value");

return oneOfType.createValue(caseType, convertedValue);
return oneOfType.createValue(caseType, converter.get(value.getValue()));
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -492,7 +492,21 @@ private boolean equivalent(Schema other, EquivalenceNullablePolicy nullablePolic

@Override
public String toString() {
return SchemaUtils.toPrettyString(this);
StringBuilder builder = new StringBuilder();
builder.append("Fields:");
builder.append(System.lineSeparator());
for (Field field : fields) {
builder.append(field);
builder.append(System.lineSeparator());
}
builder.append("Encoding positions:");
builder.append(System.lineSeparator());
builder.append(encodingPositions);
builder.append(System.lineSeparator());
builder.append("Options:");
builder.append(options);
builder.append("UUID: " + uuid);
return builder.toString();
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,21 +17,14 @@
*/
package org.apache.beam.sdk.schemas;

import java.util.Arrays;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import org.apache.beam.sdk.schemas.Schema.FieldType;
import org.apache.beam.sdk.schemas.Schema.LogicalType;
import org.apache.beam.sdk.values.Row;

/** A set of utility functions for schemas. */
@SuppressWarnings({
"nullness" // TODO(https://github.com/apache/beam/issues/20497)
})
public class SchemaUtils {
private static final String INDENT = " ";

/**
* Given two schema that have matching types, return a nullable-widened schema.
*
Expand Down Expand Up @@ -129,276 +122,4 @@ public static <BaseT, InputT> InputT toLogicalInputType(
LogicalType<InputT, BaseT> logicalType, BaseT baseType) {
return logicalType.toInputType(baseType);
}

public static String toPrettyString(Row row) {
return toPrettyRowString(row, "");
}

public static String toPrettyString(Schema schema) {
return toPrettySchemaString(schema, "");
}

static String toFieldTypeNameString(FieldType fieldType) {
return fieldType.getTypeName()
+ (Boolean.TRUE.equals(fieldType.getNullable()) ? "" : " NOT NULL");
}

static String toPrettyFieldTypeString(Schema.FieldType fieldType, String prefix) {
String nextPrefix = prefix + INDENT;
switch (fieldType.getTypeName()) {
case BYTE:
case INT16:
case INT32:
case INT64:
case DECIMAL:
case FLOAT:
case DOUBLE:
case STRING:
case DATETIME:
case BOOLEAN:
case BYTES:
return "<" + toFieldTypeNameString(fieldType) + ">";
case ARRAY:
case ITERABLE:
{
StringBuilder sb = new StringBuilder();
sb.append("<").append(toFieldTypeNameString(fieldType)).append("> {\n");
sb.append(nextPrefix)
.append("<element>: ")
.append(
toPrettyFieldTypeString(
Objects.requireNonNull(fieldType.getCollectionElementType()), nextPrefix))
.append("\n");
sb.append(prefix).append("}");
return sb.toString();
}
case MAP:
{
StringBuilder sb = new StringBuilder();
sb.append("<").append(toFieldTypeNameString(fieldType)).append("> {\n");
sb.append(nextPrefix)
.append("<key>: ")
.append(
toPrettyFieldTypeString(
Objects.requireNonNull(fieldType.getMapKeyType()), nextPrefix))
.append(",\n");
sb.append(nextPrefix)
.append("<value>: ")
.append(
toPrettyFieldTypeString(
Objects.requireNonNull(fieldType.getMapValueType()), nextPrefix))
.append("\n");
sb.append(prefix).append("}");
return sb.toString();
}
case ROW:
{
return "<"
+ toFieldTypeNameString(fieldType)
+ "> "
+ toPrettySchemaString(Objects.requireNonNull(fieldType.getRowSchema()), prefix);
}
case LOGICAL_TYPE:
{
Schema.FieldType baseType =
Objects.requireNonNull(fieldType.getLogicalType()).getBaseType();
StringBuilder sb = new StringBuilder();
sb.append("<")
.append(toFieldTypeNameString(fieldType))
.append("(")
.append(fieldType.getLogicalType().getIdentifier())
.append(")> {\n");
sb.append(nextPrefix)
.append("<base>: ")
.append(toPrettyFieldTypeString(baseType, nextPrefix))
.append("\n");
sb.append(prefix).append("}");
return sb.toString();
}
default:
throw new UnsupportedOperationException(fieldType.getTypeName() + " is not supported");
}
}

static String toPrettyOptionsString(Schema.Options options, String prefix) {
String nextPrefix = prefix + INDENT;
StringBuilder sb = new StringBuilder();
sb.append("{\n");
for (String optionName : options.getOptionNames()) {
sb.append(nextPrefix)
.append(optionName)
.append(" = ")
.append(
toPrettyFieldValueString(
options.getType(optionName), options.getValue(optionName), nextPrefix))
.append("\n");
}
sb.append(prefix).append("}");
return sb.toString();
}

static String toPrettyFieldValueString(Schema.FieldType fieldType, Object value, String prefix) {
String nextPrefix = prefix + INDENT;
switch (fieldType.getTypeName()) {
case BYTE:
case INT16:
case INT32:
case INT64:
case DECIMAL:
case FLOAT:
case DOUBLE:
case DATETIME:
case BOOLEAN:
return Objects.toString(value);
case STRING:
{
String string = (String) value;
return "\"" + string.replace("\\", "\\\\").replace("\"", "\\\"") + "\"";
}
case BYTES:
{
byte[] bytes = (byte[]) value;
return Arrays.toString(bytes);
}
case ARRAY:
case ITERABLE:
{
if (!(value instanceof List)) {
throw new IllegalArgumentException(
String.format(
"value type is '%s' for field type '%s'",
value.getClass(), fieldType.getTypeName()));
}
FieldType elementType = Objects.requireNonNull(fieldType.getCollectionElementType());

@SuppressWarnings("unchecked")
List<Object> list = (List<Object>) value;
if (list.isEmpty()) {
return "[]";
}
StringBuilder sb = new StringBuilder();
sb.append("[\n");
int size = list.size();
int index = 0;
for (Object element : list) {
sb.append(nextPrefix)
.append(toPrettyFieldValueString(elementType, element, nextPrefix));
if (index++ < size - 1) {
sb.append(",\n");
} else {
sb.append("\n");
}
}
sb.append(prefix).append("]");
return sb.toString();
}
case MAP:
{
if (!(value instanceof Map)) {
throw new IllegalArgumentException(
String.format(
"value type is '%s' for field type '%s'",
value.getClass(), fieldType.getTypeName()));
}

FieldType keyType = Objects.requireNonNull(fieldType.getMapKeyType());
FieldType valueType = Objects.requireNonNull(fieldType.getMapValueType());

@SuppressWarnings("unchecked")
Map<Object, Object> map = (Map<Object, Object>) value;
if (map.isEmpty()) {
return "{}";
}

StringBuilder sb = new StringBuilder();
sb.append("{\n");
int size = map.size();
int index = 0;
for (Map.Entry<Object, Object> entry : map.entrySet()) {
sb.append(nextPrefix)
.append(toPrettyFieldValueString(keyType, entry.getKey(), nextPrefix))
.append(": ")
.append(toPrettyFieldValueString(valueType, entry.getValue(), nextPrefix));
if (index++ < size - 1) {
sb.append(",\n");
} else {
sb.append("\n");
}
}
sb.append(prefix).append("}");
return sb.toString();
}
case ROW:
{
return toPrettyRowString((Row) value, prefix);
}
case LOGICAL_TYPE:
{
@SuppressWarnings("unchecked")
Schema.LogicalType<Object, Object> logicalType =
(Schema.LogicalType<Object, Object>)
Objects.requireNonNull(fieldType.getLogicalType());
Schema.FieldType baseType = logicalType.getBaseType();
Object baseValue = logicalType.toBaseType(value);
return toPrettyFieldValueString(baseType, baseValue, prefix);
}
default:
throw new UnsupportedOperationException(fieldType.getTypeName() + " is not supported");
}
}

static String toPrettySchemaString(Schema schema, String prefix) {
String nextPrefix = prefix + INDENT;
StringBuilder sb = new StringBuilder();
sb.append("{\n");
for (Schema.Field field : schema.getFields()) {
sb.append(nextPrefix)
.append(field.getName())
.append(": ")
.append(toPrettyFieldTypeString(field.getType(), nextPrefix));
if (field.getOptions().hasOptions()) {
sb.append(", fieldOptions = ")
.append(toPrettyOptionsString(field.getOptions(), nextPrefix));
}
sb.append("\n");
}
sb.append(prefix).append("}");
if (schema.getOptions().hasOptions()) {
sb.append(", schemaOptions = ").append(toPrettyOptionsString(schema.getOptions(), prefix));
}
if (schema.getUUID() != null) {
sb.append(", schemaUUID = ").append(schema.getUUID());
}
return sb.toString();
}

static String toPrettyRowString(Row row, String prefix) {
long nonNullFieldCount = row.getValues().stream().filter(Objects::nonNull).count();
if (nonNullFieldCount == 0) {
return "{}";
}

String nextPrefix = prefix + INDENT;
StringBuilder sb = new StringBuilder();
sb.append("{\n");
long nonNullFieldIndex = 0;
for (Schema.Field field : row.getSchema().getFields()) {
String fieldName = field.getName();
Object fieldValue = row.getValue(fieldName);
if (fieldValue == null) {
continue;
}
sb.append(nextPrefix)
.append(fieldName)
.append(": ")
.append(toPrettyFieldValueString(field.getType(), fieldValue, nextPrefix));
if (nonNullFieldIndex++ < nonNullFieldCount - 1) {
sb.append(",\n");
} else {
sb.append("\n");
}
}
sb.append(prefix).append("}");
return sb.toString();
}
}
Loading
Loading