Skip to content

Commit 668c6a3

Browse files
(IcebergIO) bugfix: propagate config properties to RecordWriter (#39250)
* (IcebergIO) bugfix: propagate config properties to RecordWriter * Refactor: add new writeProperties option to IcebergWriteSchemaTransformProvider
1 parent f8b0616 commit 668c6a3

11 files changed

Lines changed: 207 additions & 30 deletions

sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
import com.google.auto.value.AutoValue;
2323
import java.util.Arrays;
2424
import java.util.List;
25+
import java.util.Map;
2526
import org.apache.beam.sdk.annotations.Internal;
2627
import org.apache.beam.sdk.io.Read;
2728
import org.apache.beam.sdk.schemas.Schema;
@@ -406,6 +407,8 @@ public abstract static class WriteRows extends PTransform<PCollection<Row>, Iceb
406407

407408
abstract boolean getAutoSharding();
408409

410+
abstract @Nullable Map<String, String> getWriteProperties();
411+
409412
abstract Builder toBuilder();
410413

411414
@AutoValue.Builder
@@ -424,6 +427,8 @@ abstract static class Builder {
424427

425428
abstract Builder setAutoSharding(boolean autoSharding);
426429

430+
abstract Builder setWriteProperties(Map<String, String> writeProperties);
431+
427432
abstract WriteRows build();
428433
}
429434

@@ -474,6 +479,10 @@ public WriteRows withAutosharding() {
474479
return toBuilder().setAutoSharding(true).build();
475480
}
476481

482+
public WriteRows withWriteProperties(Map<String, String> writeProperties) {
483+
return toBuilder().setWriteProperties(writeProperties).build();
484+
}
485+
477486
@Override
478487
public IcebergWriteResult expand(PCollection<Row> input) {
479488
List<?> allToArgs = Arrays.asList(getTableIdentifier(), getDynamicDestinations());
@@ -509,7 +518,8 @@ public IcebergWriteResult expand(PCollection<Row> input) {
509518
getCatalogConfig(),
510519
destinations,
511520
getTriggeringFrequency(),
512-
getDirectWriteByteLimit()));
521+
getDirectWriteByteLimit(),
522+
getWriteProperties()));
513523
case HASH:
514524
return input
515525
.apply(
@@ -521,7 +531,8 @@ public IcebergWriteResult expand(PCollection<Row> input) {
521531
getCatalogConfig(),
522532
destinations,
523533
getTriggeringFrequency(),
524-
getAutoSharding()));
534+
getAutoSharding(),
535+
getWriteProperties()));
525536
default:
526537
throw new UnsupportedOperationException(
527538
"Unsupported distribution mode: " + getDistributionMode());

sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -158,6 +158,11 @@ public static Builder builder() {
158158
+ "during high-throughput writes. Only available with 'hash' distribution mode.")
159159
public abstract @Nullable Boolean getAutosharding();
160160

161+
@SchemaFieldDescription(
162+
"Properties applied to the underlying file writer (e.g. Parquet write properties like "
163+
+ "'write.parquet.bloom-filter-enabled.column.<col>').")
164+
public abstract @Nullable Map<String, String> getWriteProperties();
165+
161166
@AutoValue.Builder
162167
public abstract static class Builder {
163168
public abstract Builder setTable(String table);
@@ -188,6 +193,8 @@ public abstract static class Builder {
188193

189194
public abstract Builder setAutosharding(Boolean autosharding);
190195

196+
public abstract Builder setWriteProperties(Map<String, String> writeProperties);
197+
191198
public abstract Configuration build();
192199
}
193200

@@ -279,6 +286,11 @@ public PCollectionRowTuple expand(PCollectionRowTuple input) {
279286
writeTransform = writeTransform.withAutosharding();
280287
}
281288

289+
@Nullable Map<String, String> writeProperties = configuration.getWriteProperties();
290+
if (writeProperties != null && !writeProperties.isEmpty()) {
291+
writeTransform = writeTransform.withWriteProperties(writeProperties);
292+
}
293+
282294
// TODO: support dynamic destinations
283295
IcebergWriteResult result = rows.apply(writeTransform);
284296

sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriter.java

Lines changed: 20 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
package org.apache.beam.sdk.io.iceberg;
1919

2020
import java.io.IOException;
21+
import java.util.Map;
2122
import org.apache.beam.sdk.metrics.Counter;
2223
import org.apache.beam.sdk.metrics.Metrics;
2324
import org.apache.iceberg.DataFile;
@@ -34,6 +35,7 @@
3435
import org.apache.iceberg.io.DataWriter;
3536
import org.apache.iceberg.io.OutputFile;
3637
import org.apache.iceberg.parquet.Parquet;
38+
import org.checkerframework.checker.nullness.qual.Nullable;
3739
import org.slf4j.Logger;
3840
import org.slf4j.LoggerFactory;
3941

@@ -54,11 +56,22 @@ class RecordWriter {
5456
catalog.loadTable(destination.getTableIdentifier()),
5557
destination.getFileFormat(),
5658
filename,
57-
partitionKey);
59+
partitionKey,
60+
null);
5861
}
5962

6063
RecordWriter(Table table, FileFormat fileFormat, String filename, StructLike partitionKey)
6164
throws IOException {
65+
this(table, fileFormat, filename, partitionKey, null);
66+
}
67+
68+
RecordWriter(
69+
Table table,
70+
FileFormat fileFormat,
71+
String filename,
72+
StructLike partitionKey,
73+
@Nullable Map<String, String> writeProperties)
74+
throws IOException {
6275
this.table = table;
6376
this.fileFormat = fileFormat;
6477

@@ -91,14 +104,17 @@ class RecordWriter {
91104
.build();
92105
break;
93106
case PARQUET:
94-
icebergDataWriter =
107+
Parquet.DataWriteBuilder parquetBuilder =
95108
Parquet.writeData(outputFile)
96109
.forTable(table)
97110
.createWriterFunc(GenericParquetWriter::create)
98111
.withPartition(partitionKey)
99112
.withKeyMetadata(keyMetadata)
100-
.overwrite()
101-
.build();
113+
.overwrite();
114+
if (writeProperties != null && !writeProperties.isEmpty()) {
115+
parquetBuilder.setAll(writeProperties);
116+
}
117+
icebergDataWriter = parquetBuilder.build();
102118
break;
103119
case ORC:
104120
throw new UnsupportedOperationException("ORC file format not currently supported.");

sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -202,7 +202,8 @@ private RecordWriter createWriter(PartitionKey partitionKey) {
202202
table,
203203
icebergDestination.getFileFormat(),
204204
filePrefix + "_" + stateToken + "_" + recordIndex,
205-
partitionKey);
205+
partitionKey,
206+
writeProperties);
206207
openWriters++;
207208
return writer;
208209
} catch (IOException e) {
@@ -253,6 +254,7 @@ static String getPartitionDataPath(
253254
private final String filePrefix;
254255
private final long maxFileSize;
255256
private final int maxNumWriters;
257+
private final @Nullable Map<String, String> writeProperties;
256258
@VisibleForTesting int openWriters = 0;
257259

258260
@VisibleForTesting
@@ -265,10 +267,20 @@ static String getPartitionDataPath(
265267

266268
RecordWriterManager(
267269
IcebergCatalogConfig catalogConfig, String filePrefix, long maxFileSize, int maxNumWriters) {
270+
this(catalogConfig, filePrefix, maxFileSize, maxNumWriters, null);
271+
}
272+
273+
RecordWriterManager(
274+
IcebergCatalogConfig catalogConfig,
275+
String filePrefix,
276+
long maxFileSize,
277+
int maxNumWriters,
278+
@Nullable Map<String, String> writeProperties) {
268279
this.catalogConfig = catalogConfig;
269280
this.filePrefix = filePrefix;
270281
this.maxFileSize = maxFileSize;
271282
this.maxNumWriters = maxNumWriters;
283+
this.writeProperties = writeProperties;
272284
}
273285

274286
/**

sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteDirectRowsToFiles.java

Lines changed: 11 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -40,24 +40,27 @@ class WriteDirectRowsToFiles
4040
private final IcebergCatalogConfig catalogConfig;
4141
private final String filePrefix;
4242
private final long maxBytesPerFile;
43+
private final @Nullable Map<String, String> writeProperties;
4344

4445
WriteDirectRowsToFiles(
4546
IcebergCatalogConfig catalogConfig,
4647
DynamicDestinations dynamicDestinations,
4748
String filePrefix,
48-
long maxBytesPerFile) {
49+
long maxBytesPerFile,
50+
@Nullable Map<String, String> writeProperties) {
4951
this.catalogConfig = catalogConfig;
5052
this.dynamicDestinations = dynamicDestinations;
5153
this.filePrefix = filePrefix;
5254
this.maxBytesPerFile = maxBytesPerFile;
55+
this.writeProperties = writeProperties;
5356
}
5457

5558
@Override
5659
public PCollection<FileWriteResult> expand(PCollection<KV<String, Row>> input) {
5760
return input.apply(
5861
ParDo.of(
5962
new WriteDirectRowsToFilesDoFn(
60-
catalogConfig, dynamicDestinations, maxBytesPerFile, filePrefix)));
63+
catalogConfig, dynamicDestinations, maxBytesPerFile, filePrefix, writeProperties)));
6164
}
6265

6366
private static class WriteDirectRowsToFilesDoFn extends DoFn<KV<String, Row>, FileWriteResult> {
@@ -66,24 +69,28 @@ private static class WriteDirectRowsToFilesDoFn extends DoFn<KV<String, Row>, Fi
6669
private final IcebergCatalogConfig catalogConfig;
6770
private final String filePrefix;
6871
private final long maxFileSize;
72+
private final @Nullable Map<String, String> writeProperties;
6973
private transient @Nullable RecordWriterManager recordWriterManager;
7074

7175
WriteDirectRowsToFilesDoFn(
7276
IcebergCatalogConfig catalogConfig,
7377
DynamicDestinations dynamicDestinations,
7478
long maxFileSize,
75-
String filePrefix) {
79+
String filePrefix,
80+
@Nullable Map<String, String> writeProperties) {
7681
this.catalogConfig = catalogConfig;
7782
this.dynamicDestinations = dynamicDestinations;
7883
this.filePrefix = filePrefix;
7984
this.maxFileSize = maxFileSize;
85+
this.writeProperties = writeProperties;
8086
this.recordWriterManager = null;
8187
}
8288

8389
@StartBundle
8490
public void startBundle() {
8591
recordWriterManager =
86-
new RecordWriterManager(catalogConfig, filePrefix, maxFileSize, Integer.MAX_VALUE);
92+
new RecordWriterManager(
93+
catalogConfig, filePrefix, maxFileSize, Integer.MAX_VALUE, writeProperties);
8794
}
8895

8996
@ProcessElement

sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WriteGroupedRowsToFiles.java

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
package org.apache.beam.sdk.io.iceberg;
1919

2020
import java.util.List;
21+
import java.util.Map;
2122
import org.apache.beam.sdk.transforms.DoFn;
2223
import org.apache.beam.sdk.transforms.PTransform;
2324
import org.apache.beam.sdk.transforms.ParDo;
@@ -30,6 +31,7 @@
3031
import org.apache.beam.sdk.values.WindowedValue;
3132
import org.apache.beam.sdk.values.WindowedValues;
3233
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions;
34+
import org.checkerframework.checker.nullness.qual.Nullable;
3335

3436
class WriteGroupedRowsToFiles
3537
extends PTransform<
@@ -39,16 +41,19 @@ class WriteGroupedRowsToFiles
3941
private final DynamicDestinations dynamicDestinations;
4042
private final IcebergCatalogConfig catalogConfig;
4143
private final String filePrefix;
44+
private final @Nullable Map<String, String> writeProperties;
4245

4346
WriteGroupedRowsToFiles(
4447
IcebergCatalogConfig catalogConfig,
4548
DynamicDestinations dynamicDestinations,
4649
String filePrefix,
47-
long maxBytesPerFile) {
50+
long maxBytesPerFile,
51+
@Nullable Map<String, String> writeProperties) {
4852
this.catalogConfig = catalogConfig;
4953
this.dynamicDestinations = dynamicDestinations;
5054
this.filePrefix = filePrefix;
5155
this.maxBytesPerFile = maxBytesPerFile;
56+
this.writeProperties = writeProperties;
5257
}
5358

5459
@Override
@@ -57,7 +62,7 @@ public PCollection<FileWriteResult> expand(
5762
return input.apply(
5863
ParDo.of(
5964
new WriteGroupedRowsToFilesDoFn(
60-
catalogConfig, dynamicDestinations, maxBytesPerFile, filePrefix)));
65+
catalogConfig, dynamicDestinations, maxBytesPerFile, filePrefix, writeProperties)));
6166
}
6267

6368
private static class WriteGroupedRowsToFilesDoFn
@@ -67,16 +72,19 @@ private static class WriteGroupedRowsToFilesDoFn
6772
private final IcebergCatalogConfig catalogConfig;
6873
private final String filePrefix;
6974
private final long maxFileSize;
75+
private final @Nullable Map<String, String> writeProperties;
7076

7177
WriteGroupedRowsToFilesDoFn(
7278
IcebergCatalogConfig catalogConfig,
7379
DynamicDestinations dynamicDestinations,
7480
long maxFileSize,
75-
String filePrefix) {
81+
String filePrefix,
82+
@Nullable Map<String, String> writeProperties) {
7683
this.catalogConfig = catalogConfig;
7784
this.dynamicDestinations = dynamicDestinations;
7885
this.filePrefix = filePrefix;
7986
this.maxFileSize = maxFileSize;
87+
this.writeProperties = writeProperties;
8088
}
8189

8290
@ProcessElement
@@ -93,7 +101,8 @@ public void processElement(
93101
WindowedValues.of(destination, window.maxTimestamp(), window, paneInfo);
94102
RecordWriterManager writer;
95103
try (RecordWriterManager openWriter =
96-
new RecordWriterManager(catalogConfig, filePrefix, maxFileSize, Integer.MAX_VALUE)) {
104+
new RecordWriterManager(
105+
catalogConfig, filePrefix, maxFileSize, Integer.MAX_VALUE, writeProperties)) {
97106
writer = openWriter;
98107
for (Row e : element.getValue()) {
99108
writer.write(windowedDestination, e);

sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/WritePartitionedRowsToFiles.java

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -59,14 +59,17 @@ class WritePartitionedRowsToFiles
5959
private final DynamicDestinations dynamicDestinations;
6060
private final IcebergCatalogConfig catalogConfig;
6161
private final String filePrefix;
62+
private final @Nullable Map<String, String> writeProperties;
6263

6364
WritePartitionedRowsToFiles(
6465
IcebergCatalogConfig catalogConfig,
6566
DynamicDestinations dynamicDestinations,
66-
String filePrefix) {
67+
String filePrefix,
68+
@Nullable Map<String, String> writeProperties) {
6769
this.catalogConfig = catalogConfig;
6870
this.dynamicDestinations = dynamicDestinations;
6971
this.filePrefix = filePrefix;
72+
this.writeProperties = writeProperties;
7073
}
7174

7275
@Override
@@ -78,7 +81,9 @@ public PCollection<FileWriteResult> expand(PCollection<KV<Row, Iterable<Row>>> i
7881
.getElemCoder())
7982
.getSchema();
8083
return input.apply(
81-
ParDo.of(new WriteDoFn(catalogConfig, dynamicDestinations, filePrefix, dataSchema)));
84+
ParDo.of(
85+
new WriteDoFn(
86+
catalogConfig, dynamicDestinations, filePrefix, dataSchema, writeProperties)));
8287
}
8388

8489
private static class WriteDoFn extends DoFn<KV<Row, Iterable<Row>>, FileWriteResult> {
@@ -87,6 +92,7 @@ private static class WriteDoFn extends DoFn<KV<Row, Iterable<Row>>, FileWriteRes
8792
private final IcebergCatalogConfig catalogConfig;
8893
private final String filePrefix;
8994
private final Schema dataSchema;
95+
private final @Nullable Map<String, String> writeProperties;
9096
private transient @MonotonicNonNull Map<TableIdentifier, Integer> specIds;
9197
private transient @MonotonicNonNull Map<TableIdentifier, Map<String, PartitionField>>
9298
partitionFieldMaps;
@@ -95,11 +101,13 @@ private static class WriteDoFn extends DoFn<KV<Row, Iterable<Row>>, FileWriteRes
95101
IcebergCatalogConfig catalogConfig,
96102
DynamicDestinations dynamicDestinations,
97103
String filePrefix,
98-
Schema dataSchema) {
104+
Schema dataSchema,
105+
@Nullable Map<String, String> writeProperties) {
99106
this.catalogConfig = catalogConfig;
100107
this.dynamicDestinations = dynamicDestinations;
101108
this.filePrefix = filePrefix;
102109
this.dataSchema = dataSchema;
110+
this.writeProperties = writeProperties;
103111
}
104112

105113
@Setup
@@ -132,7 +140,8 @@ public void processElement(
132140
.addExtension(String.format("%s-%s", filePrefix, UUID.randomUUID()));
133141

134142
RecordWriter writer =
135-
new RecordWriter(table, destination.getFileFormat(), fileName, partitionData);
143+
new RecordWriter(
144+
table, destination.getFileFormat(), fileName, partitionData, writeProperties);
136145
try {
137146
for (Row row : element.getValue()) {
138147
Record record = IcebergUtils.beamRowToIcebergRecord(table.schema(), row);

0 commit comments

Comments
 (0)