apache / apache/iceberg

DynamicIcebergSink writes not being read by IcebergSource streaming

Open
#14,526 14 comments 0 reactions 0 assignees View on GitHub
bug
Dominant language
Java
Stars
9.2k
Forks
3.5k
Avg merge
2d 11h
Merged PRs (30d)
132

Description

### Apache Iceberg version

1.10.0 (latest release)

### Query engine

Flink

### Please describe the bug 🐞

Hi all,
Have 2 Flink (v1.20.3) jobs - one is writing records to an Iceberg table with org.apache.iceberg.flink.sink.FlinkSink, and another streaming incremental changes from the table with org.apache.iceberg.flink.source.IcebergSource. This was working fine.

I then changed FlinkSink to org.apache.iceberg.flink.sink.dynamic.DynamicIcebergSink. The writes to the Iceberg table are still working fine - however now the incremental reads do not get picked up by the IcebergSource in the second Flink job. The IcebergSource recognises the new snapshots, but says: "Discovered 0 splits from incremental scan". When I query the difference between the snapshots with Athena however, there definitely are incremental records there.

I couldn't find any documentation on things to consider when moving to DynamicIcebergSink - is there something that I'm missing here?

My sink is defined like this:

```java
DynamicIcebergSink.Builder flinkSinkBuilder =
DynamicIcebergSink.forInput(dataStream).generator((inputRecord, out) -> {
var record = new DynamicRecord(TableIdentifier.of(s3_table_db, s3_table_name),
"main", icebergSchema,
avroGenericRecordToRowDataMapper.map(inputRecord), partitionSpec,
partitionSpec.isPartitioned() ? DistributionMode.HASH
: DistributionMode.NONE,
2);
record.setUpsertMode(false);
out.collect(record);
}).catalogLoader(icebergCatalogLoader).immediateTableUpdate(true)
.set("write.upsert.enabled", "false");
```

Thank you.

### Willingness to contribute

- [ ] I can contribute a fix for this bug independently
- [x] I would be willing to contribute a fix for this bug with guidance from the Iceberg community
- [ ] I cannot contribute a fix for this bug at this time

Contributor guide

Open the contributing guide

Research direction

Start with org.apache.iceberg.flink.sink.dynamic.DynamicIcebergSink and org.apache.iceberg.flink.source.IcebergSource, focusing on how the sink writes snapshots and how the source performs incremental scans. Reproduce with Flink 1.20.3 and Iceberg 1.10.0, compare the discovered snapshots and splits, and confirm that records written through DynamicIcebergSink are returned by the incremental source.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
data-engineering, stream-processing
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Quiet
Clarity
Mostly clear
Newbie friendliness
48/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.