DLT Append Flow Parameterization
Options
- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
04-30-2025 01:25 AM
Hi All,
I'm currently using DLT append flow to merge multiple streaming flows into one output.
While trying to make the append flow into a dynamic function for scalability, the dlt append flow seem to have some errors.
stat_table = f"{catalog}.{bronze_schema}.output"
dlt.create_streaming_table(
name = stat_table
)
def append_flow(stat_table, source😞
@dlt.append_flow(target = f"{stat_table}", name = f"{source}_flow")
def topic_flow(topic = source😞
return(
dlt.read_stream(f"{source}")
)
list = ['table1', 'table2', 'table3']
for source in list:
append_flow(stat_table, source)
The {source} is another dlt view within the same dlt pipeline.
The error message:
Flow 'dejian.test.table3_flow' could not be planned in append mode, but there are multiple flows writing to its destination `dejian`.`bronze`.`__materialization_mat_2d2a9216_c0f8_4ca4_ad69_11fec5c94151_output_1`. Starting in complete mode will cause results to be overwritten. Please edit the flow definition to allow for append mode.
Append mode error (full traces in the driver logs):
[STREAMING_OUTPUT_MODE.UNSUPPORTED_OPERATION] Invalid streaming output mode: append. This output mode is not supported for streaming aggregations without watermark on streaming DataFrames/DataSets. SQLSTATE: 42KDE
I tried using watermarks but I think it does not work as well.
Flow 'dejian.test.table3_flow' could not be planned in append mode, but there are multiple flows writing to its destination `dejian`.`bronze`.`__materialization_mat_2d2a9216_c0f8_4ca4_ad69_11fec5c94151_output_1`. Starting in complete mode will cause results to be overwritten. Please edit the flow definition to allow for append mode.
Append mode error (full traces in the driver logs):
[STREAMING_OUTPUT_MODE.UNSUPPORTED_OPERATION] Invalid streaming output mode: append. This output mode is not supported for streaming aggregations without watermark on streaming DataFrames/DataSets. SQLSTATE: 42KDE
I tried using watermarks but I think it does not work as well.
I saw example in documentation show looping for different kafka topics as source, does this support non-kafka source as well?
Please advice.
Thank you.