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]