Job is idle — throughput ~0; structure shown.
Fragment 25343 (Actor 123609,123610)
StreamMaterialize { columns: [party_id, entity_id, entity_type, $src(hidden)], stream_key: [party_id, entity_id, $src], pk_columns: [party_id, entity_id, $src], pk_conflict: NoCheck }
├── output: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, 'account':Varchar, $src ]
├── stream key: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, $src ]
└── StreamUnion { all: true }
├── output: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, 'account':Varchar, $src ]
├── stream key: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, $src ]
├── MergeExecutor
│ ├── output: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, 'account':Varchar, 0:Int32 ]
│ └── stream key: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id ]
└── MergeExecutor
├── output: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id, 'portfolio':Varchar, 1:Int32 ]
└── stream key: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id ]
Fragment 25344 (Actor 123614,123613)
StreamProject { exprs: [active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, 'account':Varchar, 0:Int32] }
├── output: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, 'account':Varchar, 0:Int32 ]
├── stream key: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id ]
└── MergeExecutor
├── output: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, accounts_dm.account_id ]
└── stream key: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id ]
Fragment 25345 (Actor 123611,123612)
StreamSyncLogStore
├── output: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, accounts_dm.account_id ]
├── stream key: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id ]
└── StreamHashJoin { type: Inner, predicate: active_parties_to_accounts_mv_next.account_id = accounts_dm.account_id }
├── output: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, accounts_dm.account_id ]
├── stream key: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id ]
├── MergeExecutor
│ ├── output: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, active_parties_mv.id ]
│ └── stream key: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id ]
└── MergeExecutor { output: [ accounts_dm.account_id ], stream key: [ accounts_dm.account_id ] }
Fragment 25346 (Actor 123616,123615)
StreamSyncLogStore
├── output: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, active_parties_mv.id ]
├── stream key: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id ]
└── StreamHashJoin { type: Inner, predicate: active_parties_to_accounts_mv_next.party_id = active_parties_mv.id }
├── output: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id, active_parties_mv.id ]
├── stream key: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id ]
├── MergeExecutor
│ ├── output: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id ]
│ └── stream key: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id ]
└── MergeExecutor { output: [ active_parties_mv.id ], stream key: [ active_parties_mv.id ] }
Fragment 25347 (Actor 123608,123607)
StreamTableScan { table: active_parties_to_accounts_mv_next, columns: [party_id, account_id] }
├── output: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id ]
├── stream key: [ active_parties_to_accounts_mv_next.party_id, active_parties_to_accounts_mv_next.account_id ]
├── Upstream { output: [ party_id, account_id ], stream key: [] }
└── BatchPlanNode { output: [ party_id, account_id ], stream key: [] }
Fragment 25348 (Actor 123624,123623)
StreamTableScan { table: active_parties_mv, columns: [id] } { output: [ active_parties_mv.id ], stream key: [ active_parties_mv.id ] }
├── Upstream { output: [ id ], stream key: [] }
└── BatchPlanNode { output: [ id ], stream key: [] }
Fragment 25349 (Actor 123625,123626)
StreamProject { exprs: [accounts_dm.account_id] } { output: [ accounts_dm.account_id ], stream key: [ accounts_dm.account_id ] }
└── StreamFilter { predicate: IsNull(accounts_dm.disabled_at) } { output: [ accounts_dm.account_id, accounts_dm.disabled_at ], stream key: [ accounts_dm.account_id ] }
└── StreamTableScan { table: accounts_dm, columns: [account_id, disabled_at] }
├── output: [ accounts_dm.account_id, accounts_dm.disabled_at ]
├── stream key: [ accounts_dm.account_id ]
├── Upstream { output: [ account_id, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ account_id, disabled_at ], stream key: [] }
Fragment 25350 (Actor 123617,123618)
StreamProject { exprs: [active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id, 'portfolio':Varchar, 1:Int32] }
├── output: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id, 'portfolio':Varchar, 1:Int32 ]
├── stream key: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id ]
└── MergeExecutor
├── output: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id, portfolios_dm.portfolio_id ]
└── stream key: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id ]
Fragment 25351 (Actor 123620,123619)
StreamSyncLogStore
├── output: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id, portfolios_dm.portfolio_id ]
├── stream key: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id ]
└── StreamHashJoin { type: Inner, predicate: active_parties_to_portfolios_mv_next.portfolio_id = portfolios_dm.portfolio_id }
├── output: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id, portfolios_dm.portfolio_id ]
├── stream key: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id ]
├── MergeExecutor
│ ├── output: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id, active_parties_mv.id ]
│ └── stream key: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id ]
└── MergeExecutor { output: [ portfolios_dm.portfolio_id ], stream key: [ portfolios_dm.portfolio_id ] }
Fragment 25352 (Actor 123622,123621)
StreamSyncLogStore
├── output: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id, active_parties_mv.id ]
├── stream key: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id ]
└── StreamHashJoin { type: Inner, predicate: active_parties_to_portfolios_mv_next.party_id = active_parties_mv.id }
├── output: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id, active_parties_mv.id ]
├── stream key: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id ]
├── MergeExecutor
│ ├── output: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id ]
│ └── stream key: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id ]
└── MergeExecutor { output: [ active_parties_mv.id ], stream key: [ active_parties_mv.id ] }
Fragment 25353 (Actor 123627,123628)
StreamTableScan { table: active_parties_to_portfolios_mv_next, columns: [party_id, portfolio_id] }
├── output: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id ]
├── stream key: [ active_parties_to_portfolios_mv_next.party_id, active_parties_to_portfolios_mv_next.portfolio_id ]
├── Upstream { output: [ party_id, portfolio_id ], stream key: [] }
└── BatchPlanNode { output: [ party_id, portfolio_id ], stream key: [] }
Fragment 25354 (Actor 123629,123630)
StreamTableScan { table: active_parties_mv, columns: [id] } { output: [ active_parties_mv.id ], stream key: [ active_parties_mv.id ] }
├── Upstream { output: [ id ], stream key: [] }
└── BatchPlanNode { output: [ id ], stream key: [] }
Fragment 25355 (Actor 123631,123632)
StreamProject { exprs: [portfolios_dm.portfolio_id] } { output: [ portfolios_dm.portfolio_id ], stream key: [ portfolios_dm.portfolio_id ] }
└── StreamFilter { predicate: IsNull(portfolios_dm.disabled_at) } { output: [ portfolios_dm.portfolio_id, portfolios_dm.disabled_at ], stream key: [ portfolios_dm.portfolio_id ] }
└── StreamTableScan { table: portfolios_dm, columns: [portfolio_id, disabled_at] }
├── output: [ portfolios_dm.portfolio_id, portfolios_dm.disabled_at ]
├── stream key: [ portfolios_dm.portfolio_id ]
├── Upstream { output: [ portfolio_id, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ portfolio_id, disabled_at ], stream key: [] }