Job is idle — throughput ~0; structure shown.
Fragment 12337 (Actor 33581,33580)
StreamMaterialize { columns: [team_id, task_id, updated_at], stream_key: [team_id, task_id], pk_columns: [team_id, task_id], pk_conflict: NoCheck }
├── output: [ team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at ]
├── stream key: [ team_to_clients_mv.team_id, tasks_dm.task_id ]
└── StreamProject { exprs: [team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at] }
├── output: [ team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at ]
├── stream key: [ team_to_clients_mv.team_id, tasks_dm.task_id ]
└── StreamGroupTopN { order: [tasks_dm.updated_at DESC], limit: 1, offset: 0, group_key: [team_to_clients_mv.team_id, tasks_dm.task_id] }
├── output:
│ ┌── team_to_clients_mv.team_id
│ ├── tasks_dm.task_id
│ ├── tasks_dm.updated_at
│ ├── tasks_dm.task_id
│ ├── team_to_clients_mv.team_id
│ ├── tasks_dm.resource_client_id
│ └── $src
├── stream key: [ team_to_clients_mv.team_id, tasks_dm.task_id ]
└── MergeExecutor
├── output:
│ ┌── team_to_clients_mv.team_id
│ ├── tasks_dm.task_id
│ ├── tasks_dm.updated_at
│ ├── tasks_dm.task_id
│ ├── team_to_clients_mv.team_id
│ ├── tasks_dm.resource_client_id
│ └── $src
└── stream key: [ tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id, $src ]
Fragment 12338 (Actor 33583,33582)
StreamUnion { all: true }
├── output:
│ ┌── team_to_clients_mv.team_id
│ ├── tasks_dm.task_id
│ ├── tasks_dm.updated_at
│ ├── tasks_dm.task_id
│ ├── team_to_clients_mv.team_id
│ ├── tasks_dm.resource_client_id
│ └── $src
├── stream key: [ tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id, $src ]
├── MergeExecutor
│ ├── output:
│ │ ┌── team_to_clients_mv.team_id
│ │ ├── tasks_dm.task_id
│ │ ├── tasks_dm.updated_at
│ │ ├── tasks_dm.task_id
│ │ ├── team_to_clients_mv.team_id
│ │ ├── tasks_dm.resource_client_id
│ │ └── 0:Int32
│ └── stream key: [ tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id ]
├── MergeExecutor
│ ├── output:
│ │ ┌── team_to_accounts_mv_next.team_id
│ │ ├── tasks_dm.task_id
│ │ ├── tasks_dm.updated_at
│ │ ├── tasks_dm.task_id
│ │ ├── team_to_accounts_mv_next.team_id
│ │ ├── tasks_dm.resource_account_id
│ │ └── 1:Int32
│ └── stream key: [ tasks_dm.task_id, team_to_accounts_mv_next.team_id, tasks_dm.resource_account_id ]
├── MergeExecutor
│ ├── output:
│ │ ┌── team_to_portfolios_mv_next.team_id
│ │ ├── tasks_dm.task_id
│ │ ├── tasks_dm.updated_at
│ │ ├── tasks_dm.task_id
│ │ ├── team_to_portfolios_mv_next.team_id
│ │ ├── tasks_dm.resource_portfolio_id
│ │ └── 2:Int32
│ └── stream key: [ tasks_dm.task_id, team_to_portfolios_mv_next.team_id, tasks_dm.resource_portfolio_id ]
└── MergeExecutor
├── output:
│ ┌── team_to_parties_mv.team_id
│ ├── tasks_dm.task_id
│ ├── tasks_dm.updated_at
│ ├── tasks_dm.task_id
│ ├── team_to_parties_mv.team_id
│ ├── tasks_dm.resource_party_id
│ └── 3:Int32
└── stream key: [ tasks_dm.task_id, team_to_parties_mv.team_id, tasks_dm.resource_party_id ]
Fragment 12339 (Actor 33584,33585)
StreamProject { exprs: [team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id, 0:Int32] }
├── output: [ team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id, 0:Int32 ]
├── stream key: [ tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_client_id = team_to_clients_mv.client_id }
├── output: [ team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_client_id, team_to_clients_mv.client_id ]
├── stream key: [ tasks_dm.task_id, team_to_clients_mv.team_id, tasks_dm.resource_client_id ]
├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }
└── MergeExecutor { output: [ team_to_clients_mv.team_id, team_to_clients_mv.client_id ], stream key: [ team_to_clients_mv.team_id, team_to_clients_mv.client_id ] }
Fragment 12340 (Actor 33567,33566)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.task_id ]
└── StreamFilter { predicate: IsNull(tasks_dm.disabled_at) }
├── output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
├── stream key: [ tasks_dm.task_id ]
└── StreamTableScan { table: tasks_dm, columns: [task_id, resource_client_id, updated_at, disabled_at] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
├── stream key: [ tasks_dm.task_id ]
├── Upstream { output: [ task_id, resource_client_id, updated_at, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ task_id, resource_client_id, updated_at, disabled_at ], stream key: [] }
Fragment 12341 (Actor 33611,33610)
StreamTableScan { table: team_to_clients_mv, columns: [team_id, client_id] }
├── output: [ team_to_clients_mv.team_id, team_to_clients_mv.client_id ]
├── stream key: [ team_to_clients_mv.team_id, team_to_clients_mv.client_id ]
├── Upstream { output: [ team_id, client_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, client_id ], stream key: [] }
Fragment 12342 (Actor 33609,33608)
StreamProject { exprs: [team_to_accounts_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_accounts_mv_next.team_id, tasks_dm.resource_account_id, 1:Int32] }
├── output: [ team_to_accounts_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_accounts_mv_next.team_id, tasks_dm.resource_account_id, 1:Int32 ]
├── stream key: [ tasks_dm.task_id, team_to_accounts_mv_next.team_id, tasks_dm.resource_account_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_account_id = team_to_accounts_mv_next.account_id }
├── output: [ team_to_accounts_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, team_to_accounts_mv_next.account_id ]
├── stream key: [ tasks_dm.task_id, team_to_accounts_mv_next.team_id, tasks_dm.resource_account_id ]
├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }
└── MergeExecutor
├── output: [ team_to_accounts_mv_next.team_id, team_to_accounts_mv_next.account_id ]
└── stream key: [ team_to_accounts_mv_next.team_id, team_to_accounts_mv_next.account_id ]
Fragment 12343 (Actor 33568,33569)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.task_id ]
└── StreamFilter { predicate: IsNull(tasks_dm.disabled_at) }
├── output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
├── stream key: [ tasks_dm.task_id ]
└── StreamTableScan { table: tasks_dm, columns: [task_id, resource_account_id, updated_at, disabled_at] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
├── stream key: [ tasks_dm.task_id ]
├── Upstream { output: [ task_id, resource_account_id, updated_at, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ task_id, resource_account_id, updated_at, disabled_at ], stream key: [] }
Fragment 12344 (Actor 33621,33620)
StreamTableScan { table: team_to_accounts_mv_next, columns: [team_id, account_id] }
├── output: [ team_to_accounts_mv_next.team_id, team_to_accounts_mv_next.account_id ]
├── stream key: [ team_to_accounts_mv_next.team_id, team_to_accounts_mv_next.account_id ]
├── Upstream { output: [ team_id, account_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, account_id ], stream key: [] }
Fragment 12345 (Actor 33612,33613)
StreamProject { exprs: [team_to_portfolios_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_portfolios_mv_next.team_id, tasks_dm.resource_portfolio_id, 2:Int32] }
├── output: [ team_to_portfolios_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_portfolios_mv_next.team_id, tasks_dm.resource_portfolio_id, 2:Int32 ]
├── stream key: [ tasks_dm.task_id, team_to_portfolios_mv_next.team_id, tasks_dm.resource_portfolio_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_portfolio_id = team_to_portfolios_mv_next.portfolio_id }
├── output: [ team_to_portfolios_mv_next.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, team_to_portfolios_mv_next.portfolio_id ]
├── stream key: [ tasks_dm.task_id, team_to_portfolios_mv_next.team_id, tasks_dm.resource_portfolio_id ]
├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }
└── MergeExecutor
├── output: [ team_to_portfolios_mv_next.team_id, team_to_portfolios_mv_next.portfolio_id ]
└── stream key: [ team_to_portfolios_mv_next.team_id, team_to_portfolios_mv_next.portfolio_id ]
Fragment 12346 (Actor 33622,33623)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.task_id ]
└── StreamFilter { predicate: IsNull(tasks_dm.disabled_at) }
├── output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
├── stream key: [ tasks_dm.task_id ]
└── StreamTableScan { table: tasks_dm, columns: [task_id, resource_portfolio_id, updated_at, disabled_at] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
├── stream key: [ tasks_dm.task_id ]
├── Upstream { output: [ task_id, resource_portfolio_id, updated_at, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ task_id, resource_portfolio_id, updated_at, disabled_at ], stream key: [] }
Fragment 12347 (Actor 33625,33624)
StreamTableScan { table: team_to_portfolios_mv_next, columns: [team_id, portfolio_id] }
├── output: [ team_to_portfolios_mv_next.team_id, team_to_portfolios_mv_next.portfolio_id ]
├── stream key: [ team_to_portfolios_mv_next.team_id, team_to_portfolios_mv_next.portfolio_id ]
├── Upstream { output: [ team_id, portfolio_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, portfolio_id ], stream key: [] }
Fragment 12348 (Actor 33615,33614)
StreamProject { exprs: [team_to_parties_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_parties_mv.team_id, tasks_dm.resource_party_id, 3:Int32] }
├── output: [ team_to_parties_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.task_id, team_to_parties_mv.team_id, tasks_dm.resource_party_id, 3:Int32 ]
├── stream key: [ tasks_dm.task_id, team_to_parties_mv.team_id, tasks_dm.resource_party_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_party_id = team_to_parties_mv.party_id }
├── output: [ team_to_parties_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_party_id, team_to_parties_mv.party_id ]
├── stream key: [ tasks_dm.task_id, team_to_parties_mv.team_id, tasks_dm.resource_party_id ]
├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }
└── MergeExecutor { output: [ team_to_parties_mv.team_id, team_to_parties_mv.party_id ], stream key: [ team_to_parties_mv.team_id, team_to_parties_mv.party_id ] }
Fragment 12349 (Actor 33617,33616)
StreamProject { exprs: [tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.task_id ]
└── StreamFilter { predicate: IsNull(tasks_dm.disabled_at) }
├── output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
├── stream key: [ tasks_dm.task_id ]
└── StreamTableScan { table: tasks_dm, columns: [task_id, resource_party_id, updated_at, disabled_at] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at, tasks_dm.disabled_at ]
├── stream key: [ tasks_dm.task_id ]
├── Upstream { output: [ task_id, resource_party_id, updated_at, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ task_id, resource_party_id, updated_at, disabled_at ], stream key: [] }
Fragment 12350 (Actor 33627,33626)
StreamTableScan { table: team_to_parties_mv, columns: [team_id, party_id] }
├── output: [ team_to_parties_mv.team_id, team_to_parties_mv.party_id ]
├── stream key: [ team_to_parties_mv.team_id, team_to_parties_mv.party_id ]
├── Upstream { output: [ team_id, party_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, party_id ], stream key: [] }