Configurable error handling strategy when target partition does not exist during Stream Load #66686
Replies: 1 comment
|
For the 30-day retention example, filtering expired records before the Doris sink gives you a selective workaround today: route them to a Flink side output for auditing, and send only in-window records to Doris. Keep that cutoff aligned with the table's partition-retention boundary and timezone. If you only need to discard those records, the connector passes Stream Load parameters through I would avoid using That workaround doesn't cover an unexpectedly missing in-window partition or provide the proposed per-error handling strategy, but it separates intentionally expired data from genuine partition-management failures without swallowing unrelated load errors. |
Uh oh!
There was an error while loading. Please reload this page.
When using the Flink Doris Connector to write data to partitioned tables via Stream Load, if the target partition does not exist, the connector currently only has one behavior: retry and then throw a DorisBatchLoadException, which causes the entire Flink job to fail.
This is problematic in real-world streaming scenarios where:
Historical or dirty data may arrive with partition values outside the expected range (e.g., a record from last year hitting a table that only retains 30 days of partitions), causing the entire job to crash.
In multi-table sync scenarios, a single table's partition issue can bring down the entire pipeline, affecting all other tables.
Add a configurable error handling strategy for "partition does not exist" scenarios. Suggested options:
All reactions