RWM Console cluster: risingwave-adib.adib-rw.svc.cluster.local

← cluster adib_rm objects sdk_activity_transactions_merged_mv explain
Overview Objects Graph History
materialized view · adib_rm.sdk_activity_transactions_merged_mv profiled over 5s
seconds (1–30)
Stateful hash join (4 state tables) — consider a temporal join for dimension lookupsAggregation state — unbounded unless keyed or temporally filteredDynamic filter — verify it pairs with a temporal condition to clean state
165 operators
Materialize · adib_rm.sdk_activity_transactions_merged_mv
0% idle 2 actors
Union
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · IsNull(transactions_dm_next.order_id)
2 actors
Filter · IsNull(transactions_dm_next.order_id)
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or…
2 actors
HashJoin · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Project · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or…
2 actors
HashAgg · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or… Aggregation state — unbounded unless keyed or temporally filtered
5% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Filter · transactions_dm_next
0% idle 2 actors
StreamScan · transactions_dm_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Filter · IsNull(transactions_dm_next.transaction_id)
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · transactions_intraday_dm.transaction_id = transactions_dm_n…
2 actors
HashJoin · LeftOuter · transactions_intraday_dm.transaction_id = transactions_dm_n… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · transactions_dm_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Filter · LeftOuter · transactions_intraday_dm.transaction_id = transactions_dm_n…
0% idle 2 actors
Project · LeftOuter · transactions_intraday_dm.transaction_id = transactions_dm_n…
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · transactions_intraday_dm.transaction_id = transactions_intr…
2 actors
HashJoin · LeftOuter · transactions_intraday_dm.transaction_id = transactions_intr… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftSemi · transactions_intraday_dm.status_label_id = active_transacti…
2 actors
HashJoin · LeftSemi · transactions_intraday_dm.status_label_id = active_transacti… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · active_transaction_status_id_mv
0% idle 1 actor
BatchPlan
1 actor
Merge
1 actor
Merge
2 actors
Exchange
0% idle 0 actors
Project · transactions_intraday_dm
2 actors
DynamicFilter · transactions_intraday_dm Dynamic filter — verify it pairs with a temporal condition to clean state
0% idle 2 actors
Merge
2 actors
Exchange
0% 2/s 0 actors
Project
1 actor
Now
0% 2/s 1 actor
Project · transactions_intraday_dm
2 actors
Filter · transactions_intraday_dm
4% idle 2 actors
StreamScan · transactions_intraday_dm
4% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · transactions_intraday_dm
2 actors
DynamicFilter · transactions_intraday_dm Dynamic filter — verify it pairs with a temporal condition to clean state
0% idle 2 actors
Merge
2 actors
Exchange
0% 2/s 0 actors
Project
1 actor
Now
0% 2/s 1 actor
Project · transactions_intraday_dm
2 actors
Filter · transactions_intraday_dm
3% idle 2 actors
StreamScan · transactions_intraday_dm
3% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · IsNull(transactions_intraday_dm.transaction_id)
2 actors
Filter · IsNull(transactions_intraday_dm.transaction_id)
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · transactions_dm_next.transaction_id = transactions_intraday…
2 actors
HashJoin · LeftOuter · transactions_dm_next.transaction_id = transactions_intraday… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · transactions_dm_next.transaction_id = fee_transactions_dm.t…
2 actors
HashJoin · LeftOuter · transactions_dm_next.transaction_id = fee_transactions_dm.t… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · fee_transactions_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · transactions_dm_next.transaction_id = income_transactions_d…
2 actors
HashJoin · LeftOuter · transactions_dm_next.transaction_id = income_transactions_d… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · income_transactions_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · transactions_dm_next.transaction_id = trade_transactions_dm…
2 actors
HashJoin · LeftOuter · transactions_dm_next.transaction_id = trade_transactions_dm… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · trade_transactions_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · transactions_dm_next
2 actors
DynamicFilter · transactions_dm_next Dynamic filter — verify it pairs with a temporal condition to clean state
0% idle 2 actors
Merge
2 actors
Exchange
0% 2/s 0 actors
Project
1 actor
Now
0% 2/s 1 actor
Project · transactions_dm_next
2 actors
Filter · transactions_dm_next
3% idle 2 actors
StreamScan · transactions_dm_next
3% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · IsNull(transactions_dm_next.order_id)
2 actors
Filter · IsNull(transactions_dm_next.order_id)
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or…
2 actors
HashJoin · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Project · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or…
2 actors
HashAgg · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or… Aggregation state — unbounded unless keyed or temporally filtered
1% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Filter · (transactions_dm_next.order_id <> '':Varchar)
0% idle 2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
NoOp
0% idle 2 actors
Merge
2 actors
Heat = the operator's output-buffer backpressure over the sampling window. Click a node to fold its subtree.
Materialize · adib_rm.sdk_activity_transactions_merged_mv Materialize adib_rm.sdk_activity_tr… idle · 2 actors Union Union idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · IsNull(transactions_dm_next.order_id) Project IsNull(transactions_dm_… — · 2 actors Filter · IsNull(transactions_dm_next.order_id) Filter IsNull(transactions_dm_… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or… SyncLogStore LeftOuter · transaction… — · 2 actors HashJoin · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or… HashJoin LeftOuter · transaction… idle · 2 actors Project · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or… Project LeftOuter · transaction… — · 2 actors HashAgg · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or… HashAgg LeftOuter · transaction… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Filter · transactions_dm_next Filter transactions_dm_next idle · 2 actors StreamScan · transactions_dm_next StreamScan transactions_dm_next idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Filter · IsNull(transactions_dm_next.transaction_id) Filter IsNull(transactions_dm_… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · transactions_intraday_dm.transaction_id = transactions_dm_n… SyncLogStore LeftOuter · transaction… — · 2 actors HashJoin · LeftOuter · transactions_intraday_dm.transaction_id = transactions_dm_n… HashJoin LeftOuter · transaction… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · transactions_dm_next StreamScan transactions_dm_next idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Filter · LeftOuter · transactions_intraday_dm.transaction_id = transactions_dm_n… Filter LeftOuter · transaction… idle · 2 actors Project · LeftOuter · transactions_intraday_dm.transaction_id = transactions_dm_n… Project LeftOuter · transaction… — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · transactions_intraday_dm.transaction_id = transactions_intr… SyncLogStore LeftOuter · transaction… — · 2 actors HashJoin · LeftOuter · transactions_intraday_dm.transaction_id = transactions_intr… HashJoin LeftOuter · transaction… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore SyncLogStore — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore SyncLogStore — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftSemi · transactions_intraday_dm.status_label_id = active_transacti… SyncLogStore LeftSemi · transactions… — · 2 actors HashJoin · LeftSemi · transactions_intraday_dm.status_label_id = active_transacti… HashJoin LeftSemi · transactions… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · active_transaction_status_id_mv StreamScan active_transaction_stat… idle · 1 actor BatchPlan BatchPlan — · 1 actor Merge Merge — · 1 actor Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · transactions_intraday_dm Project transactions_intraday_dm — · 2 actors DynamicFilter · transactions_intraday_dm DynamicFilter transactions_intraday_dm idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Project Project — · 1 actor Now Now 2/s · 1 actor Project · transactions_intraday_dm Project transactions_intraday_dm — · 2 actors Filter · transactions_intraday_dm Filter transactions_intraday_dm idle · 2 actors StreamScan · transactions_intraday_dm StreamScan transactions_intraday_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · transactions_intraday_dm Project transactions_intraday_dm — · 2 actors DynamicFilter · transactions_intraday_dm DynamicFilter transactions_intraday_dm idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Project Project — · 1 actor Now Now 2/s · 1 actor Project · transactions_intraday_dm Project transactions_intraday_dm — · 2 actors Filter · transactions_intraday_dm Filter transactions_intraday_dm idle · 2 actors StreamScan · transactions_intraday_dm StreamScan transactions_intraday_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · IsNull(transactions_intraday_dm.transaction_id) Project IsNull(transactions_int… — · 2 actors Filter · IsNull(transactions_intraday_dm.transaction_id) Filter IsNull(transactions_int… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · transactions_dm_next.transaction_id = transactions_intraday… SyncLogStore LeftOuter · transaction… — · 2 actors HashJoin · LeftOuter · transactions_dm_next.transaction_id = transactions_intraday… HashJoin LeftOuter · transaction… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore SyncLogStore — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · transactions_dm_next.transaction_id = fee_transactions_dm.t… SyncLogStore LeftOuter · transaction… — · 2 actors HashJoin · LeftOuter · transactions_dm_next.transaction_id = fee_transactions_dm.t… HashJoin LeftOuter · transaction… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · fee_transactions_dm StreamScan fee_transactions_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore SyncLogStore — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · transactions_dm_next.transaction_id = income_transactions_d… SyncLogStore LeftOuter · transaction… — · 2 actors HashJoin · LeftOuter · transactions_dm_next.transaction_id = income_transactions_d… HashJoin LeftOuter · transaction… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · income_transactions_dm StreamScan income_transactions_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore SyncLogStore — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · transactions_dm_next.transaction_id = trade_transactions_dm… SyncLogStore LeftOuter · transaction… — · 2 actors HashJoin · LeftOuter · transactions_dm_next.transaction_id = trade_transactions_dm… HashJoin LeftOuter · transaction… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · trade_transactions_dm StreamScan trade_transactions_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · transactions_dm_next Project transactions_dm_next — · 2 actors DynamicFilter · transactions_dm_next DynamicFilter transactions_dm_next idle · 2 actors Merge Merge — · 2 actors Exchange Exchange 2/s · 0 actors Project Project — · 1 actor Now Now 2/s · 1 actor Project · transactions_dm_next Project transactions_dm_next — · 2 actors Filter · transactions_dm_next Filter transactions_dm_next idle · 2 actors StreamScan · transactions_dm_next StreamScan transactions_dm_next idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · IsNull(transactions_dm_next.order_id) Project IsNull(transactions_dm_… — · 2 actors Filter · IsNull(transactions_dm_next.order_id) Filter IsNull(transactions_dm_… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or… SyncLogStore LeftOuter · transaction… — · 2 actors HashJoin · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or… HashJoin LeftOuter · transaction… idle · 2 actors Project · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or… Project LeftOuter · transaction… — · 2 actors HashAgg · LeftOuter · transactions_intraday_dm.order_id = transactions_dm_next.or… HashAgg LeftOuter · transaction… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Filter · (transactions_dm_next.order_id <> '':Varchar) Filter (transactions_dm_next.o… idle · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors NoOp NoOp idle · 2 actors Merge Merge — · 2 actors
Streaming operator plan from EXPLAIN ANALYZE. Node heat = backpressure. Drag to pan, scroll to zoom.
Fragments (DESCRIBE FRAGMENTS) — click to expand
Fragment 19661 (Actor 96129,96128)
StreamMaterialize { columns: [transaction_id, account_id, asset_id, source_asset_id, transaction_valuation_date, transaction_valuation_timestamp, transaction_type_id, currency_code, gross_value, net_value, quantity, order_id, order_side_label_id, transaction_source, transactions_intraday_dm.transaction_id(hidden), transactions_intraday_dm.status_label_id(hidden), transactions_intraday_dm.order_id(hidden), transactions_intraday_dm.account_id(hidden), transactions_intraday_dm.asset_id(hidden), transactions_intraday_dm.transaction_type_id(hidden), $src(hidden)], stream_key: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id, $src], pk_columns: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id, $src], pk_conflict: NoCheck }
├── output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, 'INTRADAY':Varchar, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id, $src ]
├── stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id, $src ]
└── StreamUnion { all: true } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, 'INTRADAY':Varchar, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id, $src ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id, $src ] }
    ├── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, 'INTRADAY':Varchar, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id, 0:Int32 ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id ] }
    ├── MergeExecutor { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, 'EOD':Varchar, transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, 1:Int32 ], stream key: [ transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id ] }
    └── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, 'INTRADAY_SETTLING':Varchar, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id, 2:Int32 ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id ] }

Fragment 19662 (Actor 96132,96133)
StreamProject { exprs: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, 'INTRADAY':Varchar, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id, 0:Int32] }
├── output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, 'INTRADAY':Varchar, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id, 0:Int32 ]
├── stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id ]
└── StreamFilter { predicate: IsNull(transactions_dm_next.order_id) } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.order_id, transactions_intraday_dm.status_label_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id ] }
    └── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.order_id, transactions_intraday_dm.status_label_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id ] }

Fragment 19663 (Actor 96131,96130)
StreamSyncLogStore { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.order_id, transactions_intraday_dm.status_label_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id ] }
└── StreamHashJoin { type: LeftOuter, predicate: transactions_intraday_dm.order_id = transactions_dm_next.order_id AND transactions_intraday_dm.account_id = transactions_dm_next.account_id AND transactions_intraday_dm.asset_id = transactions_dm_next.asset_id AND transactions_intraday_dm.transaction_type_id = transactions_dm_next.transaction_type_id AND (max(transactions_dm_next.transaction_valuation_timestamp) >= transactions_intraday_dm.transaction_valuation_timestamp) AND (transactions_intraday_dm.order_id <> '':Varchar) }
    ├── output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.order_id, transactions_intraday_dm.status_label_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id ]
    ├── stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id ]
    ├── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
    └── StreamProject { exprs: [transactions_dm_next.order_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, max(transactions_dm_next.transaction_valuation_timestamp)] } { output: [ transactions_dm_next.order_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, max(transactions_dm_next.transaction_valuation_timestamp) ], stream key: [ transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, transactions_dm_next.order_id ] }
        └── StreamHashAgg { group_key: [transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, transactions_dm_next.order_id], aggs: [max(transactions_dm_next.transaction_valuation_timestamp), count] } { output: [ transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, transactions_dm_next.order_id, max(transactions_dm_next.transaction_valuation_timestamp), count ], stream key: [ transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, transactions_dm_next.order_id ] }
            └── MergeExecutor { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id ], stream key: [ transactions_dm_next.transaction_id ] }

Fragment 19664 (Actor 96180,96179)
StreamNoOp { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
└── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }

Fragment 19665 (Actor 96178,96177)
StreamProject { exprs: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id] }
├── output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ]
├── stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ]
└── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }

Fragment 19666 (Actor 96174,96173)
StreamSyncLogStore { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
└── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }

Fragment 19667 (Actor 96184,96183)
StreamSyncLogStore { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
└── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }

Fragment 19668 (Actor 96175,96176)
StreamSyncLogStore { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
└── StreamHashJoin { type: LeftSemi, predicate: transactions_intraday_dm.status_label_id = active_transaction_status_id_mv.status_label_id } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
    ├── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id ] }
    └── MergeExecutor { output: [ active_transaction_status_id_mv.status_label_id, active_transaction_status_id_mv._row_id ], stream key: [ active_transaction_status_id_mv._row_id ] }

Fragment 19669 (Actor 96442,96443)
StreamProject { exprs: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id] }
├── output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ]
├── stream key: [ transactions_intraday_dm.transaction_id ]
└── StreamDynamicFilter { predicate: ($expr1 >= $expr2), output: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id, $expr1] }
    ├── output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id, $expr1 ]
    ├── stream key: [ transactions_intraday_dm.transaction_id ]
    ├── StreamProject { exprs: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id, AtTimeZone(transactions_intraday_dm.transaction_valuation_date::Timestamp, 'UTC':Varchar) as $expr1] }
    │   ├── output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id, $expr1 ]
    │   ├── stream key: [ transactions_intraday_dm.transaction_id ]
    │   └── StreamFilter { predicate: IsNull(transactions_intraday_dm.disabled_at) } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.disabled_at ], stream key: [ transactions_intraday_dm.transaction_id ] }
    │       └── StreamTableScan { table: transactions_intraday_dm, columns: [transaction_id, account_id, asset_id, transaction_valuation_date, transaction_valuation_timestamp, transaction_type_id, currency_code, gross_value, net_value, quantity, order_id, order_side_label_id, status_label_id, disabled_at] } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.disabled_at ], stream key: [ transactions_intraday_dm.transaction_id ] }
    │           ├── Upstream { output: [ transaction_id, account_id, asset_id, transaction_valuation_date, transaction_valuation_timestamp, transaction_type_id, currency_code, gross_value, net_value, quantity, order_id, order_side_label_id, status_label_id, disabled_at ], stream key: [] }
    │           └── BatchPlanNode { output: [ transaction_id, account_id, asset_id, transaction_valuation_date, transaction_valuation_timestamp, transaction_type_id, currency_code, gross_value, net_value, quantity, order_id, order_side_label_id, status_label_id, disabled_at ], stream key: [] }
    └── MergeExecutor { output: [ $expr2 ], stream key: [] }

Fragment 19670 (Actor 96185)
StreamProject { exprs: [SubtractWithTimeZone(now, '2 years':Interval, 'UTC':Varchar) as $expr2] } { output: [ $expr2 ], stream key: [] }
└── StreamNow { output: [ now ], stream key: [] }

Fragment 19671 (Actor 96444)
StreamTableScan { table: active_transaction_status_id_mv, columns: [status_label_id, _row_id] } { output: [ active_transaction_status_id_mv.status_label_id, active_transaction_status_id_mv._row_id ], stream key: [ active_transaction_status_id_mv._row_id ] }
├── Upstream { output: [ status_label_id, _row_id ], stream key: [] }
└── BatchPlanNode { output: [ status_label_id, _row_id ], stream key: [] }

Fragment 19672 (Actor 96288,96289)
StreamFilter { predicate: (transactions_dm_next.order_id <> '':Varchar) } { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id ], stream key: [ transactions_dm_next.transaction_id ] }
└── MergeExecutor { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id ], stream key: [ transactions_dm_next.transaction_id ] }

Fragment 19673 (Actor 96305,96304)
StreamProject { exprs: [transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, Coalesce(income_transactions_dm.source_asset_id, fee_transactions_dm.source_asset_id) as $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id] }
├── output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id ]
├── stream key: [ transactions_dm_next.transaction_id ]
└── MergeExecutor { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, income_transactions_dm.source_asset_id, fee_transactions_dm.source_asset_id, fee_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }

Fragment 19674 (Actor 96296,96297)
StreamSyncLogStore { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, income_transactions_dm.source_asset_id, fee_transactions_dm.source_asset_id, fee_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
└── MergeExecutor { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, income_transactions_dm.source_asset_id, fee_transactions_dm.source_asset_id, fee_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }

Fragment 19675 (Actor 96300,96301)
StreamSyncLogStore { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, income_transactions_dm.source_asset_id, fee_transactions_dm.source_asset_id, fee_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
└── StreamHashJoin { type: LeftOuter, predicate: transactions_dm_next.transaction_id = fee_transactions_dm.transaction_id } { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, income_transactions_dm.source_asset_id, fee_transactions_dm.source_asset_id, fee_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
    ├── MergeExecutor { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, income_transactions_dm.source_asset_id, income_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
    └── MergeExecutor { output: [ fee_transactions_dm.transaction_id, fee_transactions_dm.source_asset_id ], stream key: [ fee_transactions_dm.transaction_id ] }

Fragment 19676 (Actor 96299,96298)
StreamSyncLogStore { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, income_transactions_dm.source_asset_id, income_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
└── MergeExecutor { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, income_transactions_dm.source_asset_id, income_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }

Fragment 19677 (Actor 96302,96303)
StreamSyncLogStore { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, income_transactions_dm.source_asset_id, income_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
└── StreamHashJoin { type: LeftOuter, predicate: transactions_dm_next.transaction_id = income_transactions_dm.transaction_id } { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, income_transactions_dm.source_asset_id, income_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
    ├── MergeExecutor { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, trade_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
    └── MergeExecutor { output: [ income_transactions_dm.transaction_id, income_transactions_dm.source_asset_id ], stream key: [ income_transactions_dm.transaction_id ] }

Fragment 19678 (Actor 96306,96307)
StreamSyncLogStore { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, trade_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
└── MergeExecutor { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, trade_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }

Fragment 19679 (Actor 96295,96294)
StreamSyncLogStore { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, trade_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
└── StreamHashJoin { type: LeftOuter, predicate: transactions_dm_next.transaction_id = trade_transactions_dm.transaction_id } { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, trade_transactions_dm.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
    ├── MergeExecutor { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id ], stream key: [ transactions_dm_next.transaction_id ] }
    └── MergeExecutor { output: [ trade_transactions_dm.transaction_id, trade_transactions_dm.order_side_label_id ], stream key: [ trade_transactions_dm.transaction_id ] }

Fragment 19680 (Actor 96313,96312)
StreamProject { exprs: [transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id] } { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id ], stream key: [ transactions_dm_next.transaction_id ] }
└── StreamDynamicFilter { predicate: ($expr3 >= $expr4), output: [transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, $expr3] } { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, $expr3 ], stream key: [ transactions_dm_next.transaction_id ] }
    ├── StreamProject { exprs: [transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, AtTimeZone(transactions_dm_next.transaction_valuation_date::Timestamp, 'UTC':Varchar) as $expr3] } { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, $expr3 ], stream key: [ transactions_dm_next.transaction_id ] }
    │   └── StreamFilter { predicate: IsNull(transactions_dm_next.disabled_at) } { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, transactions_dm_next.disabled_at ], stream key: [ transactions_dm_next.transaction_id ] }
    │       └── StreamTableScan { table: transactions_dm_next, columns: [transaction_id, account_id, asset_id, transaction_valuation_date, transaction_valuation_timestamp, transaction_type_id, currency_code, gross_value, net_value, quantity, order_id, disabled_at] } { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, transactions_dm_next.disabled_at ], stream key: [ transactions_dm_next.transaction_id ] }
    │           ├── Upstream { output: [ transaction_id, account_id, asset_id, transaction_valuation_date, transaction_valuation_timestamp, transaction_type_id, currency_code, gross_value, net_value, quantity, order_id, disabled_at ], stream key: [] }
    │           └── BatchPlanNode { output: [ transaction_id, account_id, asset_id, transaction_valuation_date, transaction_valuation_timestamp, transaction_type_id, currency_code, gross_value, net_value, quantity, order_id, disabled_at ], stream key: [] }
    └── MergeExecutor { output: [ $expr4 ], stream key: [] }

Fragment 19681 (Actor 96311)
StreamProject { exprs: [SubtractWithTimeZone(now, '2 years':Interval, 'UTC':Varchar) as $expr4] } { output: [ $expr4 ], stream key: [] }
└── StreamNow { output: [ now ], stream key: [] }

Fragment 19682 (Actor 96451,96450)
StreamTableScan { table: trade_transactions_dm, columns: [transaction_id, order_side_label_id] } { output: [ trade_transactions_dm.transaction_id, trade_transactions_dm.order_side_label_id ], stream key: [ trade_transactions_dm.transaction_id ] }
├── Upstream { output: [ transaction_id, order_side_label_id ], stream key: [] }
└── BatchPlanNode { output: [ transaction_id, order_side_label_id ], stream key: [] }

Fragment 19683 (Actor 96453,96452)
StreamTableScan { table: income_transactions_dm, columns: [transaction_id, source_asset_id] } { output: [ income_transactions_dm.transaction_id, income_transactions_dm.source_asset_id ], stream key: [ income_transactions_dm.transaction_id ] }
├── Upstream { output: [ transaction_id, source_asset_id ], stream key: [] }
└── BatchPlanNode { output: [ transaction_id, source_asset_id ], stream key: [] }

Fragment 19684 (Actor 96458,96459)
StreamTableScan { table: fee_transactions_dm, columns: [transaction_id, source_asset_id] } { output: [ fee_transactions_dm.transaction_id, fee_transactions_dm.source_asset_id ], stream key: [ fee_transactions_dm.transaction_id ] }
├── Upstream { output: [ transaction_id, source_asset_id ], stream key: [] }
└── BatchPlanNode { output: [ transaction_id, source_asset_id ], stream key: [] }

Fragment 19685 (Actor 96292,96293)
StreamProject { exprs: [transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, 'EOD':Varchar, transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, 1:Int32] }
├── output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, 'EOD':Varchar, transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id, null:Varchar, null:Varchar, null:Varchar, null:Varchar, 1:Int32 ]
├── stream key: [ transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id ]
└── StreamFilter { predicate: IsNull(transactions_intraday_dm.transaction_id) } { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id ] }
    └── MergeExecutor { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id ] }

Fragment 19686 (Actor 96290,96291)
StreamSyncLogStore { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id ] }
└── StreamHashJoin { type: LeftOuter, predicate: transactions_dm_next.transaction_id = transactions_intraday_dm.transaction_id } { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id ] }
    ├── MergeExecutor { output: [ transactions_dm_next.transaction_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, $expr5, transactions_dm_next.transaction_valuation_date, transactions_dm_next.transaction_valuation_timestamp, transactions_dm_next.transaction_type_id, transactions_dm_next.currency_code, transactions_dm_next.gross_value, transactions_dm_next.net_value, transactions_dm_next.quantity, transactions_dm_next.order_id, trade_transactions_dm.order_side_label_id ], stream key: [ transactions_dm_next.transaction_id ] }
    └── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }

Fragment 19687 (Actor 96181,96182)
StreamProject { exprs: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id] } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
└── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }

Fragment 19688 (Actor 96402,96401)
StreamProject { exprs: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, 'INTRADAY_SETTLING':Varchar, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id, 2:Int32] }
├── output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, 'INTRADAY_SETTLING':Varchar, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id, 2:Int32 ]
├── stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id ]
└── StreamFilter { predicate: IsNull(transactions_dm_next.order_id) } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.order_id, transactions_intraday_dm.status_label_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id ] }
    └── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.order_id, transactions_intraday_dm.status_label_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id ] }

Fragment 19689 (Actor 96400,96399)
StreamSyncLogStore { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.order_id, transactions_intraday_dm.status_label_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id ] }
└── StreamHashJoin { type: LeftOuter, predicate: transactions_intraday_dm.order_id = transactions_dm_next.order_id AND transactions_intraday_dm.account_id = transactions_dm_next.account_id AND transactions_intraday_dm.asset_id = transactions_dm_next.asset_id AND transactions_intraday_dm.transaction_type_id = transactions_dm_next.transaction_type_id AND (transactions_intraday_dm.order_id <> '':Varchar) }
    ├── output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.order_id, transactions_intraday_dm.status_label_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id ]
    ├── stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id, transactions_intraday_dm.order_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_type_id ]
    ├── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
    └── StreamProject { exprs: [transactions_dm_next.order_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id] } { output: [ transactions_dm_next.order_id, transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id ], stream key: [ transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, transactions_dm_next.order_id ] }
        └── StreamHashAgg { group_key: [transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, transactions_dm_next.order_id], aggs: [count] } { output: [ transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, transactions_dm_next.order_id, count ], stream key: [ transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, transactions_dm_next.order_id ] }
            └── MergeExecutor { output: [ transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, transactions_dm_next.order_id, transactions_dm_next.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }

Fragment 19690 (Actor 96418,96419)
StreamFilter { predicate: IsNull(transactions_dm_next.transaction_id) } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
└── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }

Fragment 19691 (Actor 96422,96423)
StreamSyncLogStore { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
└── StreamHashJoin { type: LeftOuter, predicate: transactions_intraday_dm.transaction_id = transactions_dm_next.transaction_id } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_dm_next.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
    ├── StreamFilter { predicate: IsNull(transactions_intraday_dm.transaction_id) } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
    │   └── StreamProject { exprs: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id] }
    │       ├── output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ]
    │       ├── stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ]
    │       └── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
    └── MergeExecutor { output: [ transactions_dm_next.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }

Fragment 19692 (Actor 96420,96421)
StreamSyncLogStore { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
└── StreamHashJoin { type: LeftOuter, predicate: transactions_intraday_dm.transaction_id = transactions_intraday_dm.transaction_id } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
    ├── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id ], stream key: [ transactions_intraday_dm.transaction_id ] }
    └── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }

Fragment 19693 (Actor 96461,96460)
StreamProject { exprs: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id] } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id ], stream key: [ transactions_intraday_dm.transaction_id ] }
└── StreamDynamicFilter { predicate: ($expr6 >= $expr7), output: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, $expr6] }
    ├── output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, $expr6 ]
    ├── stream key: [ transactions_intraday_dm.transaction_id ]
    ├── StreamProject { exprs: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, AtTimeZone(transactions_intraday_dm.transaction_valuation_date::Timestamp, 'UTC':Varchar) as $expr6] }
    │   ├── output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, $expr6 ]
    │   ├── stream key: [ transactions_intraday_dm.transaction_id ]
    │   └── StreamFilter { predicate: (transactions_intraday_dm.retire_reason = 'SETTLED':Varchar) } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.retire_reason ], stream key: [ transactions_intraday_dm.transaction_id ] }
    │       └── StreamTableScan { table: transactions_intraday_dm, columns: [transaction_id, account_id, asset_id, transaction_valuation_date, transaction_valuation_timestamp, transaction_type_id, currency_code, gross_value, net_value, quantity, order_id, order_side_label_id, retire_reason] } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.retire_reason ], stream key: [ transactions_intraday_dm.transaction_id ] }
    │           ├── Upstream { output: [ transaction_id, account_id, asset_id, transaction_valuation_date, transaction_valuation_timestamp, transaction_type_id, currency_code, gross_value, net_value, quantity, order_id, order_side_label_id, retire_reason ], stream key: [] }
    │           └── BatchPlanNode { output: [ transaction_id, account_id, asset_id, transaction_valuation_date, transaction_valuation_timestamp, transaction_type_id, currency_code, gross_value, net_value, quantity, order_id, order_side_label_id, retire_reason ], stream key: [] }
    └── MergeExecutor { output: [ $expr7 ], stream key: [] }

Fragment 19694 (Actor 96427)
StreamProject { exprs: [SubtractWithTimeZone(now, '2 years':Interval, 'UTC':Varchar) as $expr7] } { output: [ $expr7 ], stream key: [] }
└── StreamNow { output: [ now ], stream key: [] }

Fragment 19695 (Actor 96172,96171)
StreamProject { exprs: [transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id] } { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }
└── MergeExecutor { output: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.account_id, transactions_intraday_dm.asset_id, null:Varchar, transactions_intraday_dm.transaction_valuation_date, transactions_intraday_dm.transaction_valuation_timestamp, transactions_intraday_dm.transaction_type_id, transactions_intraday_dm.currency_code, transactions_intraday_dm.gross_value, transactions_intraday_dm.net_value, transactions_intraday_dm.quantity, transactions_intraday_dm.order_id, transactions_intraday_dm.order_side_label_id, transactions_intraday_dm.status_label_id ], stream key: [ transactions_intraday_dm.transaction_id, transactions_intraday_dm.status_label_id ] }

Fragment 19696 (Actor 96473,96472)
StreamTableScan { table: transactions_dm_next, columns: [transaction_id] } { output: [ transactions_dm_next.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
├── Upstream { output: [ transaction_id ], stream key: [] }
└── BatchPlanNode { output: [ transaction_id ], stream key: [] }

Fragment 19697 (Actor 96474,96475)
StreamFilter { predicate: (transactions_dm_next.order_id <> '':Varchar) } { output: [ transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, transactions_dm_next.order_id, transactions_dm_next.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
└── StreamTableScan { table: transactions_dm_next, columns: [account_id, asset_id, transaction_type_id, order_id, transaction_id] } { output: [ transactions_dm_next.account_id, transactions_dm_next.asset_id, transactions_dm_next.transaction_type_id, transactions_dm_next.order_id, transactions_dm_next.transaction_id ], stream key: [ transactions_dm_next.transaction_id ] }
    ├── Upstream { output: [ account_id, asset_id, transaction_type_id, order_id, transaction_id ], stream key: [] }
    └── BatchPlanNode { output: [ account_id, asset_id, transaction_type_id, order_id, transaction_id ], stream key: [] }