TheConen commented on issue #14330:
URL: https://github.com/apache/iceberg/issues/14330#issuecomment-5952773552

   I am still having this issue with Iceberg 1.12.0. Right now, we're forced to 
use strings instead of UUIDs in our Avro schemas.
   
   ```
   Caused by: java.lang.ClassCastException: class 
org.apache.flink.table.data.binary.BinaryStringData cannot be cast to class [B 
(org.apache.flink.table.data.binary.BinaryStringData is in unnamed module of 
loader 'app'; [B is in module java.base of loader 'bootstrap'). Failed to push 
OutputTag with id 'dynamic-forward-stream' to operator. This can occur when 
multiple OutputTags with different types but identical names are being used.
        at 
org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.pushToOperator(CopyingChainingOutput.java:81)
 ~[flink-runtime-2.1.3.jar:2.1.3]
        at 
org.apache.flink.streaming.runtime.tasks.CopyingChainingOutput.collect(CopyingChainingOutput.java:61)
 ~[flink-runtime-2.1.3.jar:2.1.3]
        at 
org.apache.flink.streaming.runtime.tasks.ChainingOutput.collectAndCheckIfChained(ChainingOutput.java:97)
 ~[flink-runtime-2.1.3.jar:2.1.3]
        at 
org.apache.flink.streaming.runtime.tasks.BroadcastingOutputCollector.collect(BroadcastingOutputCollector.java:95)
 ~[flink-runtime-2.1.3.jar:2.1.3]
        at 
org.apache.flink.streaming.api.operators.ProcessOperator$ContextImpl.output(ProcessOperator.java:101)
 ~[flink-runtime-2.1.3.jar:2.1.3]
        at 
org.apache.iceberg.flink.sink.dynamic.DynamicRecordProcessor.emit(DynamicRecordProcessor.java:213)
 ~[iceberg-flink-2.1-1.12.0.jar:?]
        at 
org.apache.iceberg.flink.sink.dynamic.DynamicRecordProcessor.collect(DynamicRecordProcessor.java:184)
 ~[iceberg-flink-2.1-1.12.0.jar:?]
        at 
org.apache.iceberg.flink.sink.dynamic.DynamicRecordProcessor.collect(DynamicRecordProcessor.java:37)
 ~[iceberg-flink-2.1-1.12.0.jar:?]
   ```
   
   We're using an Avro schema and a DynamicIcebergSink to write objects into 
Iceberg (with Nessie as the catalog).
   
   Example of an Avro schema affected by this:
   
   ```
   record RawKafkaRecord {
        uuid? id = null;
        bytes? key = null;
        bytes? value = null;
        union{null, array<Header>} headers = null;
        string? topic = null;
        int? partition = null;
        long? offset = null;
        long? timestamp = null;
        string? timestamp_type = null;
        @adjust-to-utc("true")
        timestamp_ms? ingested_at = null;
   }
   
   record Header {
        string? key = null;
        bytes? value = null;
   }
   ```
   
   POJOs are generated from this schema via the avro-maven-plugin.
   
   DynamicIcebergSink:
   
   ```
   public void append(DataStream<RawKafkaRecord> records) {
           DynamicIcebergSink.forInput(records)
               .generator(new RawKafkaRecordToDynamicRecordGenerator())
               .catalogLoader(getNessieCatalogLoader())
               .uidPrefix("dynamic-sink")
               .toBranch(nessieProperties.branch())
               .immediateTableUpdate(true)
               .append();
       }
   ```
   
   The DynamicRecordGenerator:
   ```
   public class RawKafkaRecordToDynamicRecordGenerator implements 
DynamicRecordGenerator<RawKafkaRecord> {
   
       private static final org.apache.avro.Schema AVRO_SCHEMA = 
RawKafkaRecord.getClassSchema();
   
       private static final org.apache.iceberg.Schema ICEBERG_SCHEMA = 
AvroSchemaUtil.toIceberg(AVRO_SCHEMA);
   
       @Serial
       private static final long serialVersionUID = 1L;
   
       @Override
       public void generate(RawKafkaRecord inputRecord, 
Collector<DynamicRecord> out) throws Exception {
           var genericRecord = (GenericRecord) inputRecord;
           var rowData = 
AvroGenericRecordToRowDataMapper.forAvroSchema(AVRO_SCHEMA)
               .map(genericRecord);
   
           var nessieProperties = PropertyReader.getNessieProperties();
   
           var tableIdentifier = 
TableIdentifier.of(Namespace.of(nessieProperties.namespace()),
               inputRecord.getTopic()
                   .replace(".", "-"));
           var partitionSpec = PartitionSpec.unpartitioned(); // TODO 
Partitioning
   
           var result = new DynamicRecord(tableIdentifier, 
nessieProperties.branch(), ICEBERG_SCHEMA, rowData, partitionSpec);
           out.collect(result);
       }
   }
   ```
   
   Am I doing something wrong here handling Avro and Iceberg or is this bug 
still present?


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to