RWM Console cluster: risingwave-alinma.alinma-rw.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
81 operators
Materialize · opportunity.deposit_maturity_breaches_mv
0% idle 2 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · 'FIXED':Varchar = $expr2 AND Case($expr1, ((fixed_deposit_a…
2 actors
HashJoin · Inner · 'FIXED':Varchar = $expr2 AND Case($expr1, ((fixed_deposit_a… 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_next
2 actors
Filter · opportunity_conditions_mv_next
0% idle 2 actors
StreamScan · opportunity_conditions_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
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 Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · 'FIXED':Varchar = $expr2 AND Case($expr1, ((fixed_deposit_a… SyncLogStore Inner · 'FIXED':Varchar… — · 2 actors HashJoin · Inner · 'FIXED':Varchar = $expr2 AND Case($expr1, ((fixed_deposit_a… 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_next Project opportunity_conditions_… — · 2 actors Filter · opportunity_conditions_mv_next Filter opportunity_conditions_… idle · 2 actors StreamScan · opportunity_conditions_mv_next 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 63462 (Actor 746260,746259)
StreamMaterialize { columns: [opportunity_id, resource_id, fact_date, activity_name, maturity_date, days_delta, deposit_amount, currency, 'FIXED':Varchar(hidden), opportunity_conditions_mv_next._rw_projected_row_id(hidden), opportunity_conditions_mv_next._rw_projected_row_id#1(hidden)], stream_key: ['FIXED':Varchar, resource_id, maturity_date, currency, opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1], pk_columns: ['FIXED':Varchar, resource_id, maturity_date, currency, opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1], pk_conflict: NoCheck }
├── output: [ opportunity_conditions_mv_next.opportunity_id, fixed_deposit_accounts_dm.account_id, opportunity_conditions_mv_next.as_of_date, opportunity_conditions_mv_next.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_next._rw_projected_row_id, opportunity_conditions_mv_next._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_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ]
└── StreamProject { exprs: [opportunity_conditions_mv_next.opportunity_id, fixed_deposit_accounts_dm.account_id, opportunity_conditions_mv_next.as_of_date, opportunity_conditions_mv_next.activity_name, fixed_deposit_accounts_dm.maturity_date, (fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv_next.as_of_date) as $expr3, sum(holding_values_latest_mv_next.market_value), open_accounts_mv.base_currency_code, 'FIXED':Varchar, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1] }
    ├── output: [ opportunity_conditions_mv_next.opportunity_id, fixed_deposit_accounts_dm.account_id, opportunity_conditions_mv_next.as_of_date, opportunity_conditions_mv_next.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_next._rw_projected_row_id, opportunity_conditions_mv_next._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_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._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, sum(holding_values_latest_mv_next.market_value), opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.as_of_date, 'FIXED':Varchar, $expr2, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._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_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ]

Fragment 63463 (Actor 746258,746257)
StreamSyncLogStore { output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, sum(holding_values_latest_mv_next.market_value), opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.as_of_date, 'FIXED':Varchar, $expr2, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._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_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ] }
└── StreamHashJoin { type: Inner, predicate: 'FIXED':Varchar = $expr2 AND Case($expr1, ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv_next.as_of_date) <= -1:Int32), ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv_next.as_of_date) >= 1:Int32)) AND Case((opportunity_conditions_mv_next.op = 'GTE':Varchar), ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv_next.as_of_date)::Decimal >= Case($expr1, Neg(opportunity_conditions_mv_next.threshold), opportunity_conditions_mv_next.threshold)), (opportunity_conditions_mv_next.op = 'GT':Varchar), ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv_next.as_of_date)::Decimal > Case($expr1, Neg(opportunity_conditions_mv_next.threshold), opportunity_conditions_mv_next.threshold)), (opportunity_conditions_mv_next.op = 'LTE':Varchar), ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv_next.as_of_date)::Decimal <= Case($expr1, Neg(opportunity_conditions_mv_next.threshold), opportunity_conditions_mv_next.threshold)), (opportunity_conditions_mv_next.op = 'LT':Varchar), ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv_next.as_of_date)::Decimal < Case($expr1, Neg(opportunity_conditions_mv_next.threshold), opportunity_conditions_mv_next.threshold)), (opportunity_conditions_mv_next.op = 'EQ':Varchar), ((fixed_deposit_accounts_dm.maturity_date - opportunity_conditions_mv_next.as_of_date)::Decimal = Case($expr1, Neg(opportunity_conditions_mv_next.threshold), opportunity_conditions_mv_next.threshold)), false:Boolean) }
    ├── output: [ fixed_deposit_accounts_dm.account_id, fixed_deposit_accounts_dm.maturity_date, open_accounts_mv.base_currency_code, sum(holding_values_latest_mv_next.market_value), opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.as_of_date, 'FIXED':Varchar, $expr2, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._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_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._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_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.op, opportunity_conditions_mv_next.threshold, opportunity_conditions_mv_next.as_of_date, $expr1, $expr2, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ], stream key: [ $expr2, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ] }

Fragment 63464 (Actor 746262,746261)
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 63465 (Actor 746264,746263)
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 63466 (Actor 746266,746265)
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 63467 (Actor 746268,746267)
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 63468 (Actor 746270,746269)
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 63469 (Actor 746274,746273)
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 63470 (Actor 746272,746271)
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 63471 (Actor 746243,746244)
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 63472 (Actor 746246,746245)
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 63473 (Actor 746278,746277)
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 63474 (Actor 746276,746275)
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 63475 (Actor 746284,746283)
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 63476 (Actor 746285,746286)
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 63477 (Actor 746280,746279)
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 63478 (Actor 746288,746287)
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 63479 (Actor 746281,746282)
StreamLocalityProvider { locality_columns: [$expr2] } { output: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.op, opportunity_conditions_mv_next.threshold, opportunity_conditions_mv_next.as_of_date, $expr1, $expr2, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ], stream key: [ $expr2, opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ] }
└── MergeExecutor { output: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.op, opportunity_conditions_mv_next.threshold, opportunity_conditions_mv_next.as_of_date, $expr1, $expr2, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ], stream key: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ] }

Fragment 63480 (Actor 746290,746289)
StreamProject { exprs: [opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.op, opportunity_conditions_mv_next.threshold, opportunity_conditions_mv_next.as_of_date, In(opportunity_conditions_mv_next.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_next.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_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1] } { output: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.op, opportunity_conditions_mv_next.threshold, opportunity_conditions_mv_next.as_of_date, $expr1, $expr2, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ], stream key: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ] }
└── StreamFilter { predicate: In(opportunity_conditions_mv_next.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_next.as_of_date)) } { output: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.op, opportunity_conditions_mv_next.threshold, opportunity_conditions_mv_next.as_of_date, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ], stream key: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ] }
    └── StreamTableScan { table: opportunity_conditions_mv_next, columns: [opportunity_id, activity_name, op, threshold, as_of_date, _rw_projected_row_id, _rw_projected_row_id#1] } { output: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next.activity_name, opportunity_conditions_mv_next.op, opportunity_conditions_mv_next.threshold, opportunity_conditions_mv_next.as_of_date, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._rw_projected_row_id#1 ], stream key: [ opportunity_conditions_mv_next.opportunity_id, opportunity_conditions_mv_next._rw_projected_row_id, opportunity_conditions_mv_next._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: [] }