This is an automated email from the ASF dual-hosted git repository.
rzo1 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/storm.git
The following commit(s) were added to refs/heads/master by this push:
new 853343a9f Bump iceberg from 1.11.0 to 1.12.0 (#9150)
853343a9f is described below
commit 853343a9f260323b6d0641d59bdc668805d066ed
Author: Richard Zowalla <[email protected]>
AuthorDate: Thu Oct 1 09:30:36 2026 +0200
Bump iceberg from 1.11.0 to 1.12.0 (#9150)
Iceberg 1.12 removes GenericAppenderFactory. Build the task writers on
GenericFileWriterFactory instead and port the byte-counting wrapper from
FileAppenderFactory to FileWriterFactory.
---
...Factory.java => CountingFileWriterFactory.java} | 52 +++++++---------------
.../apache/storm/iceberg/common/IcebergWriter.java | 20 ++++-----
.../iceberg/common/PartitionedRecordWriter.java | 6 +--
pom.xml | 2 +-
4 files changed, 29 insertions(+), 51 deletions(-)
diff --git
a/external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/CountingAppenderFactory.java
b/external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/CountingFileWriterFactory.java
similarity index 58%
rename from
external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/CountingAppenderFactory.java
rename to
external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/CountingFileWriterFactory.java
index 49bbda8e6..bd41ffe41 100644
---
a/external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/CountingAppenderFactory.java
+++
b/external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/CountingFileWriterFactory.java
@@ -20,22 +20,20 @@ package org.apache.storm.iceberg.common;
import java.util.ArrayList;
import java.util.List;
-import org.apache.iceberg.FileFormat;
+import org.apache.iceberg.PartitionSpec;
import org.apache.iceberg.StructLike;
import org.apache.iceberg.data.Record;
import org.apache.iceberg.deletes.EqualityDeleteWriter;
import org.apache.iceberg.deletes.PositionDeleteWriter;
import org.apache.iceberg.encryption.EncryptedOutputFile;
import org.apache.iceberg.io.DataWriter;
-import org.apache.iceberg.io.FileAppender;
-import org.apache.iceberg.io.FileAppenderFactory;
-import org.apache.iceberg.io.OutputFile;
+import org.apache.iceberg.io.FileWriterFactory;
/**
- * Wraps a {@link FileAppenderFactory} and remembers every writer it hands
out, so the state can
+ * Wraps a {@link FileWriterFactory} and remembers every data writer it hands
out, so the state can
* ask how many bytes the currently buffered window has produced.
*
- * <p>The figure is an <em>estimate</em>: {@link FileAppender#length()}
reflects what the
+ * <p>The figure is an <em>estimate</em>: {@link DataWriter#length()} reflects
what the
* underlying format has flushed, and columnar formats such as Parquet keep a
sizeable in-memory
* buffer before writing a row group. It therefore under-reports until a file
is closed, which for
* a commit threshold only means committing slightly later than the configured
size.
@@ -43,22 +41,18 @@ import org.apache.iceberg.io.OutputFile;
* <p>Closed writers are kept in the list on purpose: a rolled-over file still
counts towards the
* bytes accumulated since the last commit. {@link #reset()} drops them when
the window is flushed.
*/
-class CountingAppenderFactory implements FileAppenderFactory<Record> {
+class CountingFileWriterFactory implements FileWriterFactory<Record> {
- private final FileAppenderFactory<Record> delegate;
- private final List<FileAppender<Record>> appenders = new ArrayList<>();
+ private final FileWriterFactory<Record> delegate;
private final List<DataWriter<Record>> dataWriters = new ArrayList<>();
- CountingAppenderFactory(FileAppenderFactory<Record> delegate) {
+ CountingFileWriterFactory(FileWriterFactory<Record> delegate) {
this.delegate = delegate;
}
/** Bytes written by every writer created since the last {@link #reset()}.
*/
long estimatedBytes() {
long total = 0L;
- for (FileAppender<Record> appender : appenders) {
- total += appender.length();
- }
for (DataWriter<Record> dataWriter : dataWriters) {
total += dataWriter.length();
}
@@ -67,43 +61,27 @@ class CountingAppenderFactory implements
FileAppenderFactory<Record> {
/** Forget the writers of the window that was just committed or aborted. */
void reset() {
- appenders.clear();
dataWriters.clear();
}
@Override
- public FileAppender<Record> newAppender(OutputFile outputFile, FileFormat
format) {
- FileAppender<Record> appender = delegate.newAppender(outputFile,
format);
- appenders.add(appender);
- return appender;
- }
-
- @Override
- public FileAppender<Record> newAppender(EncryptedOutputFile outputFile,
FileFormat format) {
- FileAppender<Record> appender = delegate.newAppender(outputFile,
format);
- appenders.add(appender);
- return appender;
- }
-
- @Override
- public DataWriter<Record> newDataWriter(EncryptedOutputFile file,
FileFormat format,
+ public DataWriter<Record> newDataWriter(EncryptedOutputFile file,
PartitionSpec spec,
StructLike partition) {
- DataWriter<Record> dataWriter = delegate.newDataWriter(file, format,
partition);
+ DataWriter<Record> dataWriter = delegate.newDataWriter(file, spec,
partition);
dataWriters.add(dataWriter);
return dataWriter;
}
@Override
- public EqualityDeleteWriter<Record> newEqDeleteWriter(EncryptedOutputFile
file,
- FileFormat format,
StructLike partition) {
+ public EqualityDeleteWriter<Record>
newEqualityDeleteWriter(EncryptedOutputFile file,
+ PartitionSpec
spec, StructLike partition) {
// The sink is append-only; delete writers are never requested.
- return delegate.newEqDeleteWriter(file, format, partition);
+ return delegate.newEqualityDeleteWriter(file, spec, partition);
}
@Override
- public PositionDeleteWriter<Record> newPosDeleteWriter(EncryptedOutputFile
file,
- FileFormat format,
- StructLike
partition) {
- return delegate.newPosDeleteWriter(file, format, partition);
+ public PositionDeleteWriter<Record>
newPositionDeleteWriter(EncryptedOutputFile file,
+ PartitionSpec
spec, StructLike partition) {
+ return delegate.newPositionDeleteWriter(file, spec, partition);
}
}
diff --git
a/external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/IcebergWriter.java
b/external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/IcebergWriter.java
index 27f524434..b3f8bc317 100644
---
a/external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/IcebergWriter.java
+++
b/external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/IcebergWriter.java
@@ -32,7 +32,7 @@ import org.apache.iceberg.Table;
import org.apache.iceberg.TableProperties;
import org.apache.iceberg.catalog.Catalog;
import org.apache.iceberg.catalog.TableIdentifier;
-import org.apache.iceberg.data.GenericAppenderFactory;
+import org.apache.iceberg.data.GenericFileWriterFactory;
import org.apache.iceberg.data.Record;
import org.apache.iceberg.exceptions.AlreadyExistsException;
import org.apache.iceberg.io.OutputFileFactory;
@@ -61,7 +61,7 @@ public class IcebergWriter implements Closeable {
private Catalog catalog;
private Table table;
private TaskWriter<Record> writer;
- private CountingAppenderFactory countingAppenderFactory;
+ private CountingFileWriterFactory countingWriterFactory;
public IcebergWriter(IcebergOptions options, int taskId) {
this.options = options;
@@ -146,7 +146,7 @@ public class IcebergWriter implements Closeable {
/** Roughly how many bytes the open files hold, for size-based flushing. */
public long bufferedBytes() {
- return countingAppenderFactory == null ? 0L :
countingAppenderFactory.estimatedBytes();
+ return countingWriterFactory == null ? 0L :
countingWriterFactory.estimatedBytes();
}
/**
@@ -158,8 +158,8 @@ public class IcebergWriter implements Closeable {
}
private void resetBuffer() {
- if (countingAppenderFactory != null) {
- countingAppenderFactory.reset();
+ if (countingWriterFactory != null) {
+ countingWriterFactory.reset();
}
}
@@ -172,9 +172,9 @@ public class IcebergWriter implements Closeable {
: PropertyUtil.propertyAsLong(table.properties(),
TableProperties.WRITE_TARGET_FILE_SIZE_BYTES,
TableProperties.WRITE_TARGET_FILE_SIZE_BYTES_DEFAULT);
- CountingAppenderFactory appenderFactory =
- new CountingAppenderFactory(new GenericAppenderFactory(schema,
spec));
- this.countingAppenderFactory = appenderFactory;
+ CountingFileWriterFactory writerFactory = new
CountingFileWriterFactory(
+ new
GenericFileWriterFactory.Builder(table).dataSchema(schema).dataFileFormat(format).build());
+ this.countingWriterFactory = writerFactory;
// A fresh OutputFileFactory per file set: its random operation id
keeps file names from
// replayed tuples unique.
OutputFileFactory fileFactory = OutputFileFactory
@@ -182,9 +182,9 @@ public class IcebergWriter implements Closeable {
.format(format)
.build();
if (spec.isUnpartitioned()) {
- return new UnpartitionedWriter<>(spec, format, appenderFactory,
fileFactory, table.io(), targetFileSize);
+ return new UnpartitionedWriter<>(spec, format, writerFactory,
fileFactory, table.io(), targetFileSize);
}
- return new PartitionedRecordWriter(spec, format, appenderFactory,
fileFactory, table.io(), targetFileSize, schema);
+ return new PartitionedRecordWriter(spec, format, writerFactory,
fileFactory, table.io(), targetFileSize, schema);
}
@Override
diff --git
a/external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/PartitionedRecordWriter.java
b/external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/PartitionedRecordWriter.java
index 56e1fda8c..6cbc9315b 100644
---
a/external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/PartitionedRecordWriter.java
+++
b/external/storm-iceberg/src/main/java/org/apache/storm/iceberg/common/PartitionedRecordWriter.java
@@ -24,7 +24,7 @@ import org.apache.iceberg.PartitionSpec;
import org.apache.iceberg.Schema;
import org.apache.iceberg.data.InternalRecordWrapper;
import org.apache.iceberg.data.Record;
-import org.apache.iceberg.io.FileAppenderFactory;
+import org.apache.iceberg.io.FileWriterFactory;
import org.apache.iceberg.io.FileIO;
import org.apache.iceberg.io.OutputFileFactory;
import org.apache.iceberg.io.PartitionedFanoutWriter;
@@ -37,9 +37,9 @@ class PartitionedRecordWriter extends
PartitionedFanoutWriter<Record> {
private final PartitionKey partitionKey;
private final InternalRecordWrapper wrapper;
- PartitionedRecordWriter(PartitionSpec spec, FileFormat format,
FileAppenderFactory<Record> appenderFactory,
+ PartitionedRecordWriter(PartitionSpec spec, FileFormat format,
FileWriterFactory<Record> writerFactory,
OutputFileFactory fileFactory, FileIO io, long
targetFileSize, Schema schema) {
- super(spec, format, appenderFactory, fileFactory, io, targetFileSize);
+ super(spec, format, writerFactory, fileFactory, io, targetFileSize);
this.partitionKey = new PartitionKey(spec, schema);
this.wrapper = new InternalRecordWrapper(schema.asStruct());
}
diff --git a/pom.xml b/pom.xml
index dcf0b3e07..7136af9f5 100644
--- a/pom.xml
+++ b/pom.xml
@@ -122,7 +122,7 @@
<hbase.version>2.6.6-hadoop3</hbase.version>
<!-- Pinned deliberately: storm-iceberg must stay buildable against
the Java baseline
above, so this is bumped explicitly rather than tracking the
newest release. -->
- <iceberg.version>1.11.0</iceberg.version>
+ <iceberg.version>1.12.0</iceberg.version>
<kryo.version>5.6.2</kryo.version>
<objensis.version>3.6</objensis.version>
<jakarta.servlet.version>6.1.0</jakarta.servlet.version>