Skip to content
Open
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
1 change: 1 addition & 0 deletions CHANGES.md
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@

## New Features / Improvements

* (Java) Improved BigQuery Storage Write API `TableRow` conversion performance by using schema field ordinals as verified descriptor lookup hints ([#40188](https://github.com/apache/beam/issues/40188)).
* (Python) Expanded the SDK worker heap dump (`--experiments=enable_heap_dump`) with process RSS, CPython allocator/GC stats, and glibc `mallinfo2` native-heap/fragmentation stats to help distinguish native-heap from Python-object memory growth ([#39244](https://github.com/apache/beam/issues/39244)).

## Breaking Changes
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -683,21 +683,37 @@ public static TableFieldSchema tableFieldToProtoTableField(
}

public static class SchemaInformation {
private static final int NO_FIELD_INDEX = -1;

private final TableFieldSchema tableFieldSchema;
private final List<SchemaInformation> subFields;
private final Map<String, SchemaInformation> subFieldsByName;
private final Iterable<SchemaInformation> parentSchemas;

/**
* Position of this field in its parent schema. Descriptors built by {@link
* TableRowToStorageApiProto#getDescriptorFromTableSchema} preserve this order, so the ordinal
* is a lookup hint. The name is always verified and {@link Descriptor#findFieldByName} is the
* fallback.
*/
private final int fieldIndex;

private SchemaInformation(
TableFieldSchema tableFieldSchema, Iterable<SchemaInformation> parentSchemas) {
TableFieldSchema tableFieldSchema,
Iterable<SchemaInformation> parentSchemas,
int fieldIndex) {
this.tableFieldSchema = tableFieldSchema;
this.subFields = Lists.newArrayList();
this.subFieldsByName = Maps.newHashMap();
this.parentSchemas = parentSchemas;
this.fieldIndex = fieldIndex;
int subFieldIndex = 0;
for (TableFieldSchema field : tableFieldSchema.getFieldsList()) {
SchemaInformation schemaInformation =
new SchemaInformation(
field, Iterables.concat(this.parentSchemas, ImmutableList.of(this)));
field,
Iterables.concat(this.parentSchemas, ImmutableList.of(this)),
subFieldIndex++);
subFields.add(schemaInformation);
subFieldsByName.put(field.getName().toLowerCase(), schemaInformation);
}
Expand All @@ -708,7 +724,9 @@ private SchemaInformation(
// the new SchemaInformation is traversable upwards only.
public SchemaInformation createDescendent(TableFieldSchema tableFieldSchema) {
return new SchemaInformation(
tableFieldSchema, Iterables.concat(this.parentSchemas, ImmutableList.of(this)));
tableFieldSchema,
Iterables.concat(this.parentSchemas, ImmutableList.of(this)),
NO_FIELD_INDEX);
}

public String getFullName() {
Expand Down Expand Up @@ -762,13 +780,30 @@ public SchemaInformation getSchemaForField(int i) {
return schemaInformation;
}

private static @Nullable FieldDescriptor getFieldDescriptor(
Descriptor descriptor,
List<FieldDescriptor> descriptorFields,
@Nullable SchemaInformation fieldSchema,
String protoFieldName) {
if (fieldSchema != null
&& fieldSchema.fieldIndex != NO_FIELD_INDEX
&& fieldSchema.fieldIndex < descriptorFields.size()) {
FieldDescriptor fieldDescriptor = descriptorFields.get(fieldSchema.fieldIndex);
if (fieldDescriptor.getName().equals(protoFieldName)) {
return fieldDescriptor;
}
}
// Preserve support for callers that supply a compatible descriptor in a different order.
return descriptor.findFieldByName(protoFieldName);
}

public static SchemaInformation fromTableSchema(TableSchema tableSchema) {
TableFieldSchema root =
TableFieldSchema.newBuilder()
.addAllFields(tableSchema.getFieldsList())
.setName("root")
.build();
return new SchemaInformation(root, Collections.emptyList());
return new SchemaInformation(root, Collections.emptyList(), NO_FIELD_INDEX);
}

static SchemaInformation fromTableSchema(
Expand Down Expand Up @@ -870,6 +905,9 @@ public static Descriptor wrapDescriptorProto(DescriptorProto descriptorProto)
.map(String::toLowerCase)
.collect(toSet());
}
// Descriptor#getFields() creates an unmodifiable wrapper, so call it once per message.
List<FieldDescriptor> descriptorFields =
descriptor == null ? Collections.emptyList() : descriptor.getFields();
for (final Map.Entry<String, Object> entry : map.entrySet()) {
String key = entry.getKey().toLowerCase();
if (requiredFieldsRemaining != null) {
Expand All @@ -880,8 +918,13 @@ public static Descriptor wrapDescriptorProto(DescriptorProto descriptorProto)
BigQuerySchemaUtil.isProtoCompatible(key)
? key
: BigQuerySchemaUtil.generatePlaceholderFieldName(key);
@Nullable SchemaInformation cachedFieldSchemaInformation =
schemaInformation.subFieldsByName.get(key);
@Nullable FieldDescriptor fieldDescriptor =
(descriptor == null) ? null : descriptor.findFieldByName(protoFieldName);
descriptor == null
? null
: SchemaInformation.getFieldDescriptor(
descriptor, descriptorFields, cachedFieldSchemaInformation, protoFieldName);

if (fieldDescriptor == null) {
if (unknownFields != null) {
Expand Down Expand Up @@ -958,7 +1001,9 @@ public static Descriptor wrapDescriptorProto(DescriptorProto descriptorProto)
}

SchemaInformation fieldSchemaInformation =
schemaInformation.getSchemaForField(entry.getKey());
cachedFieldSchemaInformation == null
? schemaInformation.getSchemaForField(entry.getKey())
: cachedFieldSchemaInformation;
try {
Supplier<@Nullable TableRow> getNestedUnknown =
() -> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

import static org.apache.beam.sdk.io.gcp.bigquery.BigQueryUtils.TIMESTAMP_FORMATTER;
import static org.apache.beam.sdk.io.gcp.bigquery.TableRowToStorageApiProto.TYPE_MAP_PROTO_CONVERTERS;
import static org.junit.Assert.assertArrayEquals;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNotNull;
Expand Down Expand Up @@ -55,6 +56,10 @@
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.stream.Collectors;
import javax.annotation.Nullable;
import org.apache.beam.sdk.io.gcp.bigquery.TableRowToStorageApiProto.SchemaConversionException;
Expand Down Expand Up @@ -1903,6 +1908,147 @@ public void testNullRepeatedDescriptorFromTableSchema() throws Exception {
assertTrue(repeatednof2.isEmpty());
}

@Test
public void testMessageFromTableRowWithReorderedDescriptor() throws Exception {
TableSchema tableSchema =
new TableSchema()
.setFields(
ImmutableList.of(
new TableFieldSchema().setName("stringvalue").setType("STRING"),
new TableFieldSchema().setName("intvalue").setType("INT64")));
TableSchema reorderedTableSchema =
new TableSchema()
.setFields(
ImmutableList.of(
new TableFieldSchema().setName("intvalue").setType("INT64"),
new TableFieldSchema().setName("stringvalue").setType("STRING")));
SchemaInformation schemaInformation = SchemaInformation.fromTableSchema(tableSchema);
Descriptor reorderedDescriptor =
TableRowToStorageApiProto.getDescriptorFromTableSchema(reorderedTableSchema, true, false);

DynamicMessage message =
TableRowToStorageApiProto.messageFromTableRow(
schemaInformation,
reorderedDescriptor,
new TableRow().set("stringvalue", "string").set("intvalue", 42),
false,
false,
null,
null,
-1,
TableRowToStorageApiProto.ErrorCollector.DONT_COLLECT);

DynamicMessage expectedMessage =
DynamicMessage.newBuilder(reorderedDescriptor)
.setField(reorderedDescriptor.findFieldByName("stringvalue"), "string")
.setField(reorderedDescriptor.findFieldByName("intvalue"), 42L)
.build();
assertArrayEquals(expectedMessage.toByteArray(), message.toByteArray());
}

@Test
public void testMessageFromTableRowWithShorterDescriptor() throws Exception {
TableSchema tableSchema =
new TableSchema()
.setFields(
ImmutableList.of(
new TableFieldSchema().setName("firstvalue").setType("STRING"),
new TableFieldSchema().setName("secondvalue").setType("INT64")));
TableSchema shorterTableSchema =
new TableSchema()
.setFields(
ImmutableList.of(new TableFieldSchema().setName("secondvalue").setType("INT64")));
SchemaInformation schemaInformation = SchemaInformation.fromTableSchema(tableSchema);
Descriptor shorterDescriptor =
TableRowToStorageApiProto.getDescriptorFromTableSchema(shorterTableSchema, true, false);

DynamicMessage message =
TableRowToStorageApiProto.messageFromTableRow(
schemaInformation,
shorterDescriptor,
new TableRow().set("secondvalue", 42),
false,
false,
null,
null,
-1,
TableRowToStorageApiProto.ErrorCollector.DONT_COLLECT);

DynamicMessage expectedMessage =
DynamicMessage.newBuilder(shorterDescriptor)
.setField(shorterDescriptor.findFieldByName("secondvalue"), 42L)
.build();
assertArrayEquals(expectedMessage.toByteArray(), message.toByteArray());
}

@Test
public void testMessageFromTableRowPreservesRequiredFieldValidation() throws Exception {
TableSchema tableSchema =
new TableSchema()
.setFields(
ImmutableList.of(
new TableFieldSchema()
.setName("requiredvalue")
.setType("STRING")
.setMode("REQUIRED")));
SchemaInformation schemaInformation = SchemaInformation.fromTableSchema(tableSchema);
Descriptor descriptor =
TableRowToStorageApiProto.getDescriptorFromTableSchema(tableSchema, true, false);

thrown.expect(TableRowToStorageApiProto.SchemaMissingRequiredFieldException.class);
TableRowToStorageApiProto.messageFromTableRow(
schemaInformation,
descriptor,
new TableRow(),
false,
false,
null,
null,
-1,
TableRowToStorageApiProto.ErrorCollector.DONT_COLLECT);
}

@Test
public void testMessageFromTableRowIsThreadSafe() throws Exception {
TableSchema tableSchema =
new TableSchema()
.setFields(
ImmutableList.of(
new TableFieldSchema().setName("stringvalue").setType("STRING"),
new TableFieldSchema().setName("intvalue").setType("INT64")));
SchemaInformation schemaInformation = SchemaInformation.fromTableSchema(tableSchema);
Descriptor descriptor =
TableRowToStorageApiProto.getDescriptorFromTableSchema(tableSchema, true, false);
TableRow tableRow = new TableRow().set("stringvalue", "string").set("intvalue", 42);
byte[] expected =
DynamicMessage.newBuilder(descriptor)
.setField(descriptor.findFieldByName("stringvalue"), "string")
.setField(descriptor.findFieldByName("intvalue"), 42L)
.build()
.toByteArray();
Callable<byte[]> conversion =
() ->
TableRowToStorageApiProto.messageFromTableRow(
schemaInformation,
descriptor,
tableRow,
false,
false,
null,
null,
-1,
TableRowToStorageApiProto.ErrorCollector.DONT_COLLECT)
.toByteArray();
ExecutorService executor = Executors.newFixedThreadPool(4);
try {
for (Future<byte[]> result : executor.invokeAll(Collections.nCopies(100, conversion))) {
assertArrayEquals(expected, result.get());
}
} finally {
executor.shutdownNow();
}
}

@Test
public void testIntegerTypeConversion() throws DescriptorValidationException {
String intFieldName = "int_field";
Expand Down
Loading