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>

Reply via email to