Skip to content
Merged
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 @@ -22,6 +22,7 @@
import com.google.auto.value.AutoValue;
import java.util.Arrays;
import java.util.List;
import java.util.Map;
import org.apache.beam.sdk.annotations.Internal;
import org.apache.beam.sdk.io.Read;
import org.apache.beam.sdk.schemas.Schema;
Expand Down Expand Up @@ -406,6 +407,8 @@ public abstract static class WriteRows extends PTransform<PCollection<Row>, Iceb

abstract boolean getAutoSharding();

abstract @Nullable Map<String, String> getWriteProperties();

abstract Builder toBuilder();

@AutoValue.Builder
Expand All @@ -424,6 +427,8 @@ abstract static class Builder {

abstract Builder setAutoSharding(boolean autoSharding);

abstract Builder setWriteProperties(Map<String, String> writeProperties);

abstract WriteRows build();
}

Expand Down Expand Up @@ -474,6 +479,10 @@ public WriteRows withAutosharding() {
return toBuilder().setAutoSharding(true).build();
}

public WriteRows withWriteProperties(Map<String, String> writeProperties) {
return toBuilder().setWriteProperties(writeProperties).build();
}

@Override
public IcebergWriteResult expand(PCollection<Row> input) {
List<?> allToArgs = Arrays.asList(getTableIdentifier(), getDynamicDestinations());
Expand Down Expand Up @@ -509,7 +518,8 @@ public IcebergWriteResult expand(PCollection<Row> input) {
getCatalogConfig(),
destinations,
getTriggeringFrequency(),
getDirectWriteByteLimit()));
getDirectWriteByteLimit(),
getWriteProperties()));
case HASH:
return input
.apply(
Expand All @@ -521,7 +531,8 @@ public IcebergWriteResult expand(PCollection<Row> input) {
getCatalogConfig(),
destinations,
getTriggeringFrequency(),
getAutoSharding()));
getAutoSharding(),
getWriteProperties()));
default:
throw new UnsupportedOperationException(
"Unsupported distribution mode: " + getDistributionMode());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -158,6 +158,11 @@ public static Builder builder() {
+ "during high-throughput writes. Only available with 'hash' distribution mode.")
public abstract @Nullable Boolean getAutosharding();

@SchemaFieldDescription(
"Properties applied to the underlying file writer (e.g. Parquet write properties like "
+ "'write.parquet.bloom-filter-enabled.column.<col>').")
public abstract @Nullable Map<String, String> getWriteProperties();

@AutoValue.Builder
public abstract static class Builder {
public abstract Builder setTable(String table);
Expand Down Expand Up @@ -188,6 +193,8 @@ public abstract static class Builder {

public abstract Builder setAutosharding(Boolean autosharding);

public abstract Builder setWriteProperties(Map<String, String> writeProperties);

public abstract Configuration build();
}

Expand Down Expand Up @@ -279,6 +286,11 @@ public PCollectionRowTuple expand(PCollectionRowTuple input) {
writeTransform = writeTransform.withAutosharding();
}

@Nullable Map<String, String> writeProperties = configuration.getWriteProperties();
if (writeProperties != null && !writeProperties.isEmpty()) {
writeTransform = writeTransform.withWriteProperties(writeProperties);
}

// TODO: support dynamic destinations
IcebergWriteResult result = rows.apply(writeTransform);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
package org.apache.beam.sdk.io.iceberg;

import java.io.IOException;
import java.util.Map;
import org.apache.beam.sdk.metrics.Counter;
import org.apache.beam.sdk.metrics.Metrics;
import org.apache.iceberg.DataFile;
Expand All @@ -34,6 +35,7 @@
import org.apache.iceberg.io.DataWriter;
import org.apache.iceberg.io.OutputFile;
import org.apache.iceberg.parquet.Parquet;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand All @@ -54,11 +56,22 @@ class RecordWriter {
catalog.loadTable(destination.getTableIdentifier()),
destination.getFileFormat(),
filename,
partitionKey);
partitionKey,
null);
}

RecordWriter(Table table, FileFormat fileFormat, String filename, StructLike partitionKey)
throws IOException {
this(table, fileFormat, filename, partitionKey, null);
}

RecordWriter(
Table table,
FileFormat fileFormat,
String filename,
StructLike partitionKey,
@Nullable Map<String, String> writeProperties)
throws IOException {
this.table = table;
this.fileFormat = fileFormat;

Expand Down Expand Up @@ -91,14 +104,17 @@ class RecordWriter {
.build();
break;
case PARQUET:
icebergDataWriter =
Parquet.DataWriteBuilder parquetBuilder =
Parquet.writeData(outputFile)
.forTable(table)
.createWriterFunc(GenericParquetWriter::create)
.withPartition(partitionKey)
.withKeyMetadata(keyMetadata)
.overwrite()
.build();
.overwrite();
if (writeProperties != null && !writeProperties.isEmpty()) {
parquetBuilder.setAll(writeProperties);
}
icebergDataWriter = parquetBuilder.build();
break;
case ORC:
throw new UnsupportedOperationException("ORC file format not currently supported.");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -202,7 +202,8 @@ private RecordWriter createWriter(PartitionKey partitionKey) {
table,
icebergDestination.getFileFormat(),
filePrefix + "_" + stateToken + "_" + recordIndex,
partitionKey);
partitionKey,
writeProperties);
openWriters++;
return writer;
} catch (IOException e) {
Expand Down Expand Up @@ -253,6 +254,7 @@ static String getPartitionDataPath(
private final String filePrefix;
private final long maxFileSize;
private final int maxNumWriters;
private final @Nullable Map<String, String> writeProperties;
@VisibleForTesting int openWriters = 0;

@VisibleForTesting
Expand All @@ -265,10 +267,20 @@ static String getPartitionDataPath(

RecordWriterManager(
IcebergCatalogConfig catalogConfig, String filePrefix, long maxFileSize, int maxNumWriters) {
this(catalogConfig, filePrefix, maxFileSize, maxNumWriters, null);
}

RecordWriterManager(
IcebergCatalogConfig catalogConfig,
String filePrefix,
long maxFileSize,
int maxNumWriters,
@Nullable Map<String, String> writeProperties) {
this.catalogConfig = catalogConfig;
this.filePrefix = filePrefix;
this.maxFileSize = maxFileSize;
this.maxNumWriters = maxNumWriters;
this.writeProperties = writeProperties;
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,24 +40,27 @@ class WriteDirectRowsToFiles
private final IcebergCatalogConfig catalogConfig;
private final String filePrefix;
private final long maxBytesPerFile;
private final @Nullable Map<String, String> writeProperties;

WriteDirectRowsToFiles(
IcebergCatalogConfig catalogConfig,
DynamicDestinations dynamicDestinations,
String filePrefix,
long maxBytesPerFile) {
long maxBytesPerFile,
@Nullable Map<String, String> writeProperties) {
this.catalogConfig = catalogConfig;
this.dynamicDestinations = dynamicDestinations;
this.filePrefix = filePrefix;
this.maxBytesPerFile = maxBytesPerFile;
this.writeProperties = writeProperties;
}

@Override
public PCollection<FileWriteResult> expand(PCollection<KV<String, Row>> input) {
return input.apply(
ParDo.of(
new WriteDirectRowsToFilesDoFn(
catalogConfig, dynamicDestinations, maxBytesPerFile, filePrefix)));
catalogConfig, dynamicDestinations, maxBytesPerFile, filePrefix, writeProperties)));
}

private static class WriteDirectRowsToFilesDoFn extends DoFn<KV<String, Row>, FileWriteResult> {
Expand All @@ -66,24 +69,28 @@ private static class WriteDirectRowsToFilesDoFn extends DoFn<KV<String, Row>, Fi
private final IcebergCatalogConfig catalogConfig;
private final String filePrefix;
private final long maxFileSize;
private final @Nullable Map<String, String> writeProperties;
private transient @Nullable RecordWriterManager recordWriterManager;

WriteDirectRowsToFilesDoFn(
IcebergCatalogConfig catalogConfig,
DynamicDestinations dynamicDestinations,
long maxFileSize,
String filePrefix) {
String filePrefix,
@Nullable Map<String, String> writeProperties) {
this.catalogConfig = catalogConfig;
this.dynamicDestinations = dynamicDestinations;
this.filePrefix = filePrefix;
this.maxFileSize = maxFileSize;
this.writeProperties = writeProperties;
this.recordWriterManager = null;
}

@StartBundle
public void startBundle() {
recordWriterManager =
new RecordWriterManager(catalogConfig, filePrefix, maxFileSize, Integer.MAX_VALUE);
new RecordWriterManager(
catalogConfig, filePrefix, maxFileSize, Integer.MAX_VALUE, writeProperties);
}

@ProcessElement
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
package org.apache.beam.sdk.io.iceberg;

import java.util.List;
import java.util.Map;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.PTransform;
import org.apache.beam.sdk.transforms.ParDo;
Expand All @@ -30,6 +31,7 @@
import org.apache.beam.sdk.values.WindowedValue;
import org.apache.beam.sdk.values.WindowedValues;
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions;
import org.checkerframework.checker.nullness.qual.Nullable;

class WriteGroupedRowsToFiles
extends PTransform<
Expand All @@ -39,16 +41,19 @@ class WriteGroupedRowsToFiles
private final DynamicDestinations dynamicDestinations;
private final IcebergCatalogConfig catalogConfig;
private final String filePrefix;
private final @Nullable Map<String, String> writeProperties;

WriteGroupedRowsToFiles(
IcebergCatalogConfig catalogConfig,
DynamicDestinations dynamicDestinations,
String filePrefix,
long maxBytesPerFile) {
long maxBytesPerFile,
@Nullable Map<String, String> writeProperties) {
this.catalogConfig = catalogConfig;
this.dynamicDestinations = dynamicDestinations;
this.filePrefix = filePrefix;
this.maxBytesPerFile = maxBytesPerFile;
this.writeProperties = writeProperties;
}

@Override
Expand All @@ -57,7 +62,7 @@ public PCollection<FileWriteResult> expand(
return input.apply(
ParDo.of(
new WriteGroupedRowsToFilesDoFn(
catalogConfig, dynamicDestinations, maxBytesPerFile, filePrefix)));
catalogConfig, dynamicDestinations, maxBytesPerFile, filePrefix, writeProperties)));
}

private static class WriteGroupedRowsToFilesDoFn
Expand All @@ -67,16 +72,19 @@ private static class WriteGroupedRowsToFilesDoFn
private final IcebergCatalogConfig catalogConfig;
private final String filePrefix;
private final long maxFileSize;
private final @Nullable Map<String, String> writeProperties;

WriteGroupedRowsToFilesDoFn(
IcebergCatalogConfig catalogConfig,
DynamicDestinations dynamicDestinations,
long maxFileSize,
String filePrefix) {
String filePrefix,
@Nullable Map<String, String> writeProperties) {
this.catalogConfig = catalogConfig;
this.dynamicDestinations = dynamicDestinations;
this.filePrefix = filePrefix;
this.maxFileSize = maxFileSize;
this.writeProperties = writeProperties;
}

@ProcessElement
Expand All @@ -93,7 +101,8 @@ public void processElement(
WindowedValues.of(destination, window.maxTimestamp(), window, paneInfo);
RecordWriterManager writer;
try (RecordWriterManager openWriter =
new RecordWriterManager(catalogConfig, filePrefix, maxFileSize, Integer.MAX_VALUE)) {
new RecordWriterManager(
catalogConfig, filePrefix, maxFileSize, Integer.MAX_VALUE, writeProperties)) {
writer = openWriter;
for (Row e : element.getValue()) {
writer.write(windowedDestination, e);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,14 +59,17 @@ class WritePartitionedRowsToFiles
private final DynamicDestinations dynamicDestinations;
private final IcebergCatalogConfig catalogConfig;
private final String filePrefix;
private final @Nullable Map<String, String> writeProperties;

WritePartitionedRowsToFiles(
IcebergCatalogConfig catalogConfig,
DynamicDestinations dynamicDestinations,
String filePrefix) {
String filePrefix,
@Nullable Map<String, String> writeProperties) {
this.catalogConfig = catalogConfig;
this.dynamicDestinations = dynamicDestinations;
this.filePrefix = filePrefix;
this.writeProperties = writeProperties;
}

@Override
Expand All @@ -78,7 +81,9 @@ public PCollection<FileWriteResult> expand(PCollection<KV<Row, Iterable<Row>>> i
.getElemCoder())
.getSchema();
return input.apply(
ParDo.of(new WriteDoFn(catalogConfig, dynamicDestinations, filePrefix, dataSchema)));
ParDo.of(
new WriteDoFn(
catalogConfig, dynamicDestinations, filePrefix, dataSchema, writeProperties)));
}

private static class WriteDoFn extends DoFn<KV<Row, Iterable<Row>>, FileWriteResult> {
Expand All @@ -87,6 +92,7 @@ private static class WriteDoFn extends DoFn<KV<Row, Iterable<Row>>, FileWriteRes
private final IcebergCatalogConfig catalogConfig;
private final String filePrefix;
private final Schema dataSchema;
private final @Nullable Map<String, String> writeProperties;
private transient @MonotonicNonNull Map<TableIdentifier, Integer> specIds;
private transient @MonotonicNonNull Map<TableIdentifier, Map<String, PartitionField>>
partitionFieldMaps;
Expand All @@ -95,11 +101,13 @@ private static class WriteDoFn extends DoFn<KV<Row, Iterable<Row>>, FileWriteRes
IcebergCatalogConfig catalogConfig,
DynamicDestinations dynamicDestinations,
String filePrefix,
Schema dataSchema) {
Schema dataSchema,
@Nullable Map<String, String> writeProperties) {
this.catalogConfig = catalogConfig;
this.dynamicDestinations = dynamicDestinations;
this.filePrefix = filePrefix;
this.dataSchema = dataSchema;
this.writeProperties = writeProperties;
}

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

RecordWriter writer =
new RecordWriter(table, destination.getFileFormat(), fileName, partitionData);
new RecordWriter(
table, destination.getFileFormat(), fileName, partitionData, writeProperties);
try {
for (Row row : element.getValue()) {
Record record = IcebergUtils.beamRowToIcebergRecord(table.schema(), row);
Expand Down
Loading
Loading