RWM Console cluster: risingwave-ai-hub-int.ai-hub-rwm.svc.cluster.local

← cluster opportunity objects deposit_maturity_breaches_mv explain
Overview Objects Graph History
materialized view · opportunity.deposit_maturity_breaches_mv profiled over 5s
seconds (1–30)

Job is idle — throughput ~0; structure shown.

Stateful hash join (4 state tables) — consider a temporal join for dimension lookupsAggregation state — unbounded unless keyed or temporally filtered
82 operators
Materialize · opportunity.deposit_maturity_breaches_mv
0% idle 2 actors
Project · Case($expr1, ((fixed_deposit_accounts_dm.maturity_date - op…
2 actors
Filter · Case($expr1, ((fixed_deposit_accounts_dm.maturity_date - op…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · 'FIXED':Varchar = $expr2
2 actors
HashJoin · Inner · 'FIXED':Varchar = $expr2 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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · opportunity_conditions_mv
2 actors
Filter · opportunity_conditions_mv
0% idle 2 actors
StreamScan · opportunity_conditions_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
HashAgg Aggregation state — unbounded unless keyed or temporally filtered
0% idle 2 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · LeftOuter · fixed_deposit_accounts_dm.account_id = holding_values_lates…
2 actors
HashJoin · LeftOuter · fixed_deposit_accounts_dm.account_id = holding_values_lates… 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
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Filter · holding_values_latest_mv_next
0% idle 2 actors
StreamScan · holding_values_latest_mv_next
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Union
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · structured_deposit_accounts_dm.account_id = open_accounts_m…
2 actors
HashJoin · Inner · structured_deposit_accounts_dm.account_id = open_accounts_m… 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 · open_accounts_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · structured_deposit_accounts_dm
2 actors
Filter · structured_deposit_accounts_dm
0% idle 2 actors
StreamScan · structured_deposit_accounts_dm
0% idle 2 actors
BatchPlan
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 · Inner · fixed_deposit_accounts_dm.account_id = open_accounts_mv.acc…
2 actors
HashJoin · Inner · fixed_deposit_accounts_dm.account_id = open_accounts_mv.acc… 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 · open_accounts_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · fixed_deposit_accounts_dm
2 actors
Filter · fixed_deposit_accounts_dm
0% idle 2 actors
StreamScan · fixed_deposit_accounts_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Heat = the operator's output-buffer backpressure over the sampling window. Click a node to fold its subtree.
Materialize · opportunity.deposit_maturity_breaches_mv Materialize opportunity.deposit_mat… idle · 2 actors Project · Case($expr1, ((fixed_deposit_accounts_dm.maturity_date - op… Project Case($expr1, ((fixed_de… — · 2 actors Filter · Case($expr1, ((fixed_deposit_accounts_dm.maturity_date - op… Filter Case($expr1, ((fixed_de… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · 'FIXED':Varchar = $expr2 SyncLogStore Inner · 'FIXED':Varchar… — · 2 actors HashJoin · Inner · 'FIXED':Varchar = $expr2 HashJoin Inner · 'FIXED':Varchar… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · opportunity_conditions_mv Project opportunity_conditions_… — · 2 actors Filter · opportunity_conditions_mv Filter opportunity_conditions_… idle · 2 actors StreamScan · opportunity_conditions_mv StreamScan opportunity_conditions_… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors HashAgg HashAgg idle · 2 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · LeftOuter · fixed_deposit_accounts_dm.account_id = holding_values_lates… SyncLogStore LeftOuter · fixed_depos… — · 2 actors HashJoin · LeftOuter · fixed_deposit_accounts_dm.account_id = holding_values_lates… HashJoin LeftOuter · fixed_depos… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Filter · holding_values_latest_mv_next Filter holding_values_latest_m… idle · 2 actors StreamScan · holding_values_latest_mv_next StreamScan holding_values_latest_m… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Union Union 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 SyncLogStore · Inner · structured_deposit_accounts_dm.account_id = open_accounts_m… SyncLogStore Inner · structured_depo… — · 2 actors HashJoin · Inner · structured_deposit_accounts_dm.account_id = open_accounts_m… HashJoin Inner · structured_depo… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · open_accounts_mv StreamScan open_accounts_mv idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · structured_deposit_accounts_dm Project structured_deposit_acco… — · 2 actors Filter · structured_deposit_accounts_dm Filter structured_deposit_acco… idle · 2 actors StreamScan · structured_deposit_accounts_dm StreamScan structured_deposit_acco… idle · 2 actors BatchPlan BatchPlan — · 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 · Inner · fixed_deposit_accounts_dm.account_id = open_accounts_mv.acc… SyncLogStore Inner · fixed_deposit_a… — · 2 actors HashJoin · Inner · fixed_deposit_accounts_dm.account_id = open_accounts_mv.acc… HashJoin Inner · fixed_deposit_a… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · open_accounts_mv StreamScan open_accounts_mv idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · fixed_deposit_accounts_dm Project fixed_deposit_accounts_… — · 2 actors Filter · fixed_deposit_accounts_dm Filter fixed_deposit_accounts_… idle · 2 actors StreamScan · fixed_deposit_accounts_dm StreamScan fixed_deposit_accounts_… idle · 2 actors BatchPlan BatchPlan — · 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 19559 (Actor 164850,164851)
StreamMaterialize { columns: [opportunity_id, resource_id, fact_date, activity_name, maturity_date, days_delta, deposit_amount, currency, 'FIXED':Varchar(hidden), opportunity_conditions_mv._rw_projected_row_id(hidden), opportunity_conditions_mv._rw_projected_row_id#1(hidden)], stream_key: ['FIXED':Varchar, resource_id, maturity_date, currency, opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1], pk_columns: ['FIXED':Varchar, resource_id, maturity_date, currency, opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1], pk_conflict: NoCheck } { output: [ opportunity_conditions_mv.opportunity_id, fixed_deposit_accounts_dm.account_id, opportunity_conditions_mv.as_of_date, opportunity_conditions_mv.activity_name, fixed_deposit_accounts_dm.maturity_date, $expr3, sum(holding_values_latest_mv_next.market_value), open_accounts_mv.base_currency_code, 'FIXED':Varchar, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ], stream key: [ 'FIXED':Varchar, fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ] }
└── StreamProject { exprs: [opportunity_conditions_mv.opportunity_id, fixed_deposit_accounts_dm.account_id, opportunity_conditions_mv.as_of_date, opportunity_conditions_mv.activity_name, fixed_deposit_accounts_dm.maturity_date, (fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv.as_of_date) as $expr3, sum(holding_values_latest_mv_next.market_value), open_accounts_mv.base_currency_code, 'FIXED':Varchar, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1] } { output: [ opportunity_conditions_mv.opportunity_id, fixed_deposit_accounts_dm.account_id, opportunity_conditions_mv.as_of_date, opportunity_conditions_mv.activity_name, fixed_deposit_accounts_dm.maturity_date, $expr3, sum(holding_values_latest_mv_next.market_value), open_accounts_mv.base_currency_code, 'FIXED':Varchar, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ], stream key: [ 'FIXED':Varchar, fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ] }
    └── StreamFilter { predicate: Case($expr1, ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv.as_of_date) <= -1:Int32), ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv.as_of_date) >= 1:Int32)) AND Case((opportunity_conditions_mv.op = 'GTE':Varchar), ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv.as_of_date)::Decimal >= Case($expr1, Neg(opportunity_conditions_mv.threshold), opportunity_conditions_mv.threshold)), (opportunity_conditions_mv.op = 'GT':Varchar), ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv.as_of_date)::Decimal > Case($expr1, Neg(opportunity_conditions_mv.threshold), opportunity_conditions_mv.threshold)), (opportunity_conditions_mv.op = 'LTE':Varchar), ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv.as_of_date)::Decimal <= Case($expr1, Neg(opportunity_conditions_mv.threshold), opportunity_conditions_mv.threshold)), (opportunity_conditions_mv.op = 'LT':Varchar), ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv.as_of_date)::Decimal < Case($expr1, Neg(opportunity_conditions_mv.threshold), opportunity_conditions_mv.threshold)), (opportunity_conditions_mv.op = 'EQ':Varchar), ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv.as_of_date)::Decimal = Case($expr1, Neg(opportunity_conditions_mv.threshold), opportunity_conditions_mv.threshold)), false:Boolean) }
        ├── output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, sum(holding_values_latest_mv_next.market_value), opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv.activity_name, opportunity_conditions_mv.op, opportunity_conditions_mv.threshold, opportunity_conditions_mv.as_of_date, $expr1, $expr2, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ]
        ├── stream key: [ 'FIXED':Varchar, fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ]
        └── MergeExecutor { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, sum(holding_values_latest_mv_next.market_value), opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv.activity_name, opportunity_conditions_mv.op, opportunity_conditions_mv.threshold, opportunity_conditions_mv.as_of_date, $expr1, $expr2, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ], stream key: [ 'FIXED':Varchar, fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ] }

Fragment 19560 (Actor 164882,164883)
StreamSyncLogStore { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, sum(holding_values_latest_mv_next.market_value), opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv.activity_name, opportunity_conditions_mv.op, opportunity_conditions_mv.threshold, opportunity_conditions_mv.as_of_date, $expr1, $expr2, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ], stream key: [ 'FIXED':Varchar, fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ] }
└── StreamHashJoin { type: Inner, predicate: 'FIXED':Varchar = $expr2 } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, sum(holding_values_latest_mv_next.market_value), opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv.activity_name, opportunity_conditions_mv.op, opportunity_conditions_mv.threshold, opportunity_conditions_mv.as_of_date, $expr1, $expr2, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ], stream key: [ 'FIXED':Varchar, fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ] }
    ├── MergeExecutor { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, sum(holding_values_latest_mv_next.market_value) ], stream key: [ 'FIXED':Varchar, fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code ] }
    └── MergeExecutor { output: [ opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv.activity_name, opportunity_conditions_mv.op, opportunity_conditions_mv.threshold, opportunity_conditions_mv.as_of_date, $expr1, $expr2, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ], stream key: [ $expr2, opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ] }

Fragment 19561 (Actor 164880,164881)
StreamLocalityProvider { locality_columns: ['FIXED':Varchar] } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, sum(holding_values_latest_mv_next.market_value) ], stream key: [ 'FIXED':Varchar, fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code ] }
└── MergeExecutor { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, sum(holding_values_latest_mv_next.market_value) ], stream key: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar ] }

Fragment 19562 (Actor 164871,164870)
StreamProject { exprs: [fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, sum(holding_values_latest_mv_next.market_value)] } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, sum(holding_values_latest_mv_next.market_value) ], stream key: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar ] }
└── StreamHashAgg { group_key: [fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar], aggs: [sum(holding_values_latest_mv_next.market_value), count] } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, sum(holding_values_latest_mv_next.market_value), count ], stream key: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar ] }
    └── StreamLocalityProvider { locality_columns: [fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar] } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, holding_values_latest_mv_next.market_value, $src, holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ], stream key: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, $src, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ] }
        └── MergeExecutor { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, holding_values_latest_mv_next.market_value, $src, holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ], stream key: [ fixed_deposit_accounts_dm.account_id, $src, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ] }

Fragment 19563 (Actor 164856,164857)
StreamSyncLogStore { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, holding_values_latest_mv_next.market_value, $src, holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ], stream key: [ fixed_deposit_accounts_dm.account_id, $src, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ] }
└── StreamHashJoin { type: LeftOuter, predicate: fixed_deposit_accounts_dm.account_id = holding_values_latest_mv_next.account_id } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, holding_values_latest_mv_next.market_value, $src, holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ], stream key: [ fixed_deposit_accounts_dm.account_id, $src, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ] }
    ├── MergeExecutor { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, $src ], stream key: [ fixed_deposit_accounts_dm.account_id, $src ] }
    └── MergeExecutor { output: [ holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.market_value, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ], stream key: [ holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ] }

Fragment 19564 (Actor 164855,164854)
StreamLocalityProvider { locality_columns: [fixed_deposit_accounts_dm.account_id] } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, $src ], stream key: [ fixed_deposit_accounts_dm.account_id, $src ] }
└── MergeExecutor { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, $src ], stream key: [ fixed_deposit_accounts_dm.account_id, $src ] }

Fragment 19565 (Actor 164866,164867)
StreamUnion { all: true } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, $src ], stream key: [ fixed_deposit_accounts_dm.account_id, $src ] }
├── MergeExecutor { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, 0:Int32 ], stream key: [ fixed_deposit_accounts_dm.account_id ] }
└── MergeExecutor { output: [ structured_deposit_accounts_dm.account_id, structured_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'STRUCTURED':Varchar, 1:Int32 ], stream key: [ structured_deposit_accounts_dm.account_id ] }

Fragment 19566 (Actor 164852,164853)
StreamProject { exprs: [fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, 0:Int32] } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'FIXED':Varchar, 0:Int32 ], stream key: [ fixed_deposit_accounts_dm.account_id ] }
└── MergeExecutor { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, open_accounts_mv.account_id ], stream key: [ fixed_deposit_accounts_dm.account_id ] }

Fragment 19567 (Actor 164859,164858)
StreamSyncLogStore { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, open_accounts_mv.account_id ], stream key: [ fixed_deposit_accounts_dm.account_id ] }
└── StreamHashJoin { type: Inner, predicate: fixed_deposit_accounts_dm.account_id = open_accounts_mv.account_id } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, open_accounts_mv.account_id ], stream key: [ fixed_deposit_accounts_dm.account_id ] }
    ├── MergeExecutor { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date ], stream key: [ fixed_deposit_accounts_dm.account_id ] }
    └── MergeExecutor { output: [ open_accounts_mv.account_id, open_accounts_mv.base_currency_code ], stream key: [ open_accounts_mv.account_id ] }

Fragment 19568 (Actor 164848,164849)
StreamProject { exprs: [fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date] } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date ], stream key: [ fixed_deposit_accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(fixed_deposit_accounts_dm.disabled_at) } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, fixed_deposit_accounts_dm.disabled_at ], stream key: [ fixed_deposit_accounts_dm.account_id ] }
    └── StreamTableScan { table: fixed_deposit_accounts_dm, columns: [account_id, maturity_date, disabled_at] } { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, fixed_deposit_accounts_dm.disabled_at ], stream key: [ fixed_deposit_accounts_dm.account_id ] }
        ├── Upstream { output: [ account_id, maturity_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, maturity_date, disabled_at ], stream key: [] }

Fragment 19569 (Actor 164860,164861)
StreamTableScan { table: open_accounts_mv, columns: [account_id, base_currency_code] } { output: [ open_accounts_mv.account_id, open_accounts_mv.base_currency_code ], stream key: [ open_accounts_mv.account_id ] }
├── Upstream { output: [ account_id, base_currency_code ], stream key: [] }
└── BatchPlanNode { output: [ account_id, base_currency_code ], stream key: [] }

Fragment 19570 (Actor 164879,164878)
StreamProject { exprs: [structured_deposit_accounts_dm.account_id, structured_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'STRUCTURED':Varchar, 1:Int32] } { output: [ structured_deposit_accounts_dm.account_id, structured_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, 'STRUCTURED':Varchar, 1:Int32 ], stream key: [ structured_deposit_accounts_dm.account_id ] }
└── MergeExecutor { output: [ structured_deposit_accounts_dm.account_id, structured_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, open_accounts_mv.account_id ], stream key: [ structured_deposit_accounts_dm.account_id ] }

Fragment 19571 (Actor 164872,164873)
StreamSyncLogStore { output: [ structured_deposit_accounts_dm.account_id, structured_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, open_accounts_mv.account_id ], stream key: [ structured_deposit_accounts_dm.account_id ] }
└── StreamHashJoin { type: Inner, predicate: structured_deposit_accounts_dm.account_id = open_accounts_mv.account_id } { output: [ structured_deposit_accounts_dm.account_id, structured_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, open_accounts_mv.account_id ], stream key: [ structured_deposit_accounts_dm.account_id ] }
    ├── MergeExecutor { output: [ structured_deposit_accounts_dm.account_id, structured_deposit_accounts_dm.maturity_date ], stream key: [ structured_deposit_accounts_dm.account_id ] }
    └── MergeExecutor { output: [ open_accounts_mv.account_id, open_accounts_mv.base_currency_code ], stream key: [ open_accounts_mv.account_id ] }

Fragment 19572 (Actor 164868,164869)
StreamProject { exprs: [structured_deposit_accounts_dm.account_id, structured_deposit_accounts_dm.maturity_date] } { output: [ structured_deposit_accounts_dm.account_id, structured_deposit_accounts_dm.maturity_date ], stream key: [ structured_deposit_accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(structured_deposit_accounts_dm.disabled_at) } { output: [ structured_deposit_accounts_dm.account_id, structured_deposit_accounts_dm.maturity_date, structured_deposit_accounts_dm.disabled_at ], stream key: [ structured_deposit_accounts_dm.account_id ] }
    └── StreamTableScan { table: structured_deposit_accounts_dm, columns: [account_id, maturity_date, disabled_at] } { output: [ structured_deposit_accounts_dm.account_id, structured_deposit_accounts_dm.maturity_date, structured_deposit_accounts_dm.disabled_at ], stream key: [ structured_deposit_accounts_dm.account_id ] }
        ├── Upstream { output: [ account_id, maturity_date, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ account_id, maturity_date, disabled_at ], stream key: [] }

Fragment 19573 (Actor 164875,164874)
StreamTableScan { table: open_accounts_mv, columns: [account_id, base_currency_code] } { output: [ open_accounts_mv.account_id, open_accounts_mv.base_currency_code ], stream key: [ open_accounts_mv.account_id ] }
├── Upstream { output: [ account_id, base_currency_code ], stream key: [] }
└── BatchPlanNode { output: [ account_id, base_currency_code ], stream key: [] }

Fragment 19574 (Actor 164884,164885)
StreamLocalityProvider { locality_columns: [holding_values_latest_mv_next.account_id] } { output: [ holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.market_value, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ], stream key: [ holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ] }
└── MergeExecutor { output: [ holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.market_value, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ], stream key: [ holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ] }

Fragment 19575 (Actor 164862,164863)
StreamFilter { predicate: (holding_values_latest_mv_next.type = 'ASSET':Varchar) } { output: [ holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.market_value, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ], stream key: [ holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ] }
└── StreamTableScan { table: holding_values_latest_mv_next, columns: [account_id, market_value, asset_id, type] } { output: [ holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.market_value, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ], stream key: [ holding_values_latest_mv_next.account_id, holding_values_latest_mv_next.asset_id, holding_values_latest_mv_next.type ] }
    ├── Upstream { output: [ account_id, market_value, asset_id, type ], stream key: [] }
    └── BatchPlanNode { output: [ account_id, market_value, asset_id, type ], stream key: [] }

Fragment 19576 (Actor 164864,164865)
StreamLocalityProvider { locality_columns: [$expr2] } { output: [ opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv.activity_name, opportunity_conditions_mv.op, opportunity_conditions_mv.threshold, opportunity_conditions_mv.as_of_date, $expr1, $expr2, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ], stream key: [ $expr2, opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ] }
└── MergeExecutor { output: [ opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv.activity_name, opportunity_conditions_mv.op, opportunity_conditions_mv.threshold, opportunity_conditions_mv.as_of_date, $expr1, $expr2, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ], stream key: [ opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ] }

Fragment 19577 (Actor 164876,164877)
StreamProject { exprs: [opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv.activity_name, opportunity_conditions_mv.op, opportunity_conditions_mv.threshold, opportunity_conditions_mv.as_of_date, In(opportunity_conditions_mv.activity_name, 'GET_DAYS_PAST_FIXED_DEPOSIT_ACCOUNT_MATURITY_DATE':Varchar, 'GET_DAYS_PAST_STRUCTURED_DEPOSIT_ACCOUNT_MATURITY_DATE':Varchar) as $expr1, Case(In(opportunity_conditions_mv.activity_name, 'GET_DAYS_TO_FIXED_DEPOSIT_ACCOUNT_MATURITY_DATE':Varchar, 'GET_DAYS_PAST_FIXED_DEPOSIT_ACCOUNT_MATURITY_DATE':Varchar), 'FIXED':Varchar, 'STRUCTURED':Varchar) as $expr2, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1] } { output: [ opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv.activity_name, opportunity_conditions_mv.op, opportunity_conditions_mv.threshold, opportunity_conditions_mv.as_of_date, $expr1, $expr2, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ], stream key: [ opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ] }
└── StreamFilter { predicate: In(opportunity_conditions_mv.activity_name, 'GET_DAYS_TO_FIXED_DEPOSIT_ACCOUNT_MATURITY_DATE':Varchar, 'GET_DAYS_PAST_FIXED_DEPOSIT_ACCOUNT_MATURITY_DATE':Varchar, 'GET_DAYS_TO_STRUCTURED_DEPOSIT_ACCOUNT_MATURITY_DATE':Varchar, 'GET_DAYS_PAST_STRUCTURED_DEPOSIT_ACCOUNT_MATURITY_DATE':Varchar) AND Not(IsNull(opportunity_conditions_mv.as_of_date)) } { output: [ opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv.activity_name, opportunity_conditions_mv.op, opportunity_conditions_mv.threshold, opportunity_conditions_mv.as_of_date, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ], stream key: [ opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ] }
    └── StreamTableScan { table: opportunity_conditions_mv, columns: [opportunity_id, activity_name, op, threshold, as_of_date, _rw_projected_row_id, _rw_projected_row_id#1] } { output: [ opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv.activity_name, opportunity_conditions_mv.op, opportunity_conditions_mv.threshold, opportunity_conditions_mv.as_of_date, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ], stream key: [ opportunity_conditions_mv.opportunity_id, opportunity_conditions_mv._rw_projected_row_id, opportunity_conditions_mv._rw_projected_row_id#1 ] }
        ├── Upstream { output: [ opportunity_id, activity_name, op, threshold, as_of_date, _rw_projected_row_id, _rw_projected_row_id#1 ], stream key: [] }
        └── BatchPlanNode { output: [ opportunity_id, activity_name, op, threshold, as_of_date, _rw_projected_row_id, _rw_projected_row_id#1 ], stream key: [] }