Job is idle — throughput ~0; structure shown.
Fragment 52676 (Actor 738170,738169)
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.resource_client_id
│ ├── tasks_dm.task_id
│ ├── team_to_clients_mv.team_id
│ └── $src
├── stream key: [ team_to_clients_mv.team_id, tasks_dm.task_id ]
└── StreamLocalityProvider { locality_columns: [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.resource_client_id
│ ├── tasks_dm.task_id
│ ├── team_to_clients_mv.team_id
│ └── $src
├── stream key:
│ ┌── team_to_clients_mv.team_id
│ ├── tasks_dm.task_id
│ ├── tasks_dm.resource_client_id
│ ├── tasks_dm.task_id
│ ├── team_to_clients_mv.team_id
│ └── $src
└── MergeExecutor
├── output:
│ ┌── team_to_clients_mv.team_id
│ ├── tasks_dm.task_id
│ ├── tasks_dm.updated_at
│ ├── tasks_dm.resource_client_id
│ ├── tasks_dm.task_id
│ ├── team_to_clients_mv.team_id
│ └── $src
└── stream key: [ tasks_dm.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id, $src ]
Fragment 52677 (Actor 738171,738172)
StreamUnion { all: true }
├── output:
│ ┌── team_to_clients_mv.team_id
│ ├── tasks_dm.task_id
│ ├── tasks_dm.updated_at
│ ├── tasks_dm.resource_client_id
│ ├── tasks_dm.task_id
│ ├── team_to_clients_mv.team_id
│ └── $src
├── stream key: [ tasks_dm.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id, $src ]
├── MergeExecutor
│ ├── output:
│ │ ┌── team_to_clients_mv.team_id
│ │ ├── tasks_dm.task_id
│ │ ├── tasks_dm.updated_at
│ │ ├── tasks_dm.resource_client_id
│ │ ├── tasks_dm.task_id
│ │ ├── team_to_clients_mv.team_id
│ │ └── 0:Int32
│ └── stream key: [ tasks_dm.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id ]
├── MergeExecutor
│ ├── output:
│ │ ┌── team_to_accounts_mv.team_id
│ │ ├── tasks_dm.task_id
│ │ ├── tasks_dm.updated_at
│ │ ├── tasks_dm.resource_account_id
│ │ ├── tasks_dm.task_id
│ │ ├── team_to_accounts_mv.team_id
│ │ └── 1:Int32
│ └── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id ]
├── MergeExecutor
│ ├── output:
│ │ ┌── team_to_portfolios_mv.team_id
│ │ ├── tasks_dm.task_id
│ │ ├── tasks_dm.updated_at
│ │ ├── tasks_dm.resource_portfolio_id
│ │ ├── tasks_dm.task_id
│ │ ├── team_to_portfolios_mv.team_id
│ │ └── 2:Int32
│ └── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id ]
└── MergeExecutor
├── output:
│ ┌── team_to_parties_mv.team_id
│ ├── tasks_dm.task_id
│ ├── tasks_dm.updated_at
│ ├── tasks_dm.resource_party_id
│ ├── tasks_dm.task_id
│ ├── team_to_parties_mv.team_id
│ └── 3:Int32
└── stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_id ]
Fragment 52678 (Actor 738174,738173)
StreamProject { exprs: [team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id, 0:Int32] }
├── output: [ team_to_clients_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id, 0:Int32 ]
├── stream key: [ tasks_dm.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id ]
└── MergeExecutor
├── 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.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id ]
Fragment 52679 (Actor 738175,738176)
StreamSyncLogStore
├── 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.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_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.resource_client_id, tasks_dm.task_id, team_to_clients_mv.team_id ]
├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at ], stream key: [ tasks_dm.resource_client_id, tasks_dm.task_id ] }
└── MergeExecutor { output: [ team_to_clients_mv.team_id, team_to_clients_mv.client_id ], stream key: [ team_to_clients_mv.client_id, team_to_clients_mv.team_id ] }
Fragment 52680 (Actor 738178,738177)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_client_id] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.resource_client_id, tasks_dm.task_id ]
└── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_client_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }
Fragment 52681 (Actor 738228,738227)
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 52682 (Actor 738179,738180)
StreamLocalityProvider { locality_columns: [team_to_clients_mv.client_id] }
├── output: [ team_to_clients_mv.team_id, team_to_clients_mv.client_id ]
├── stream key: [ team_to_clients_mv.client_id, team_to_clients_mv.team_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 52683 (Actor 738230,738229)
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 52684 (Actor 738185,738186)
StreamProject { exprs: [team_to_accounts_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id, 1:Int32] }
├── output: [ team_to_accounts_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id, 1:Int32 ]
├── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id ]
└── MergeExecutor
├── output: [ team_to_accounts_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, team_to_accounts_mv.account_id ]
└── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id ]
Fragment 52685 (Actor 738184,738183)
StreamSyncLogStore
├── output: [ team_to_accounts_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, team_to_accounts_mv.account_id ]
├── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_account_id = team_to_accounts_mv.account_id }
├── output: [ team_to_accounts_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_account_id, team_to_accounts_mv.account_id ]
├── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id, team_to_accounts_mv.team_id ]
├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at ], stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id ] }
└── MergeExecutor { output: [ team_to_accounts_mv.team_id, team_to_accounts_mv.account_id ], stream key: [ team_to_accounts_mv.account_id, team_to_accounts_mv.team_id ] }
Fragment 52686 (Actor 738188,738187)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_account_id] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id ]
└── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_account_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }
Fragment 52687 (Actor 738232,738231)
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 52688 (Actor 738192,738191)
StreamLocalityProvider { locality_columns: [team_to_accounts_mv.account_id] }
├── output: [ team_to_accounts_mv.team_id, team_to_accounts_mv.account_id ]
├── stream key: [ team_to_accounts_mv.account_id, team_to_accounts_mv.team_id ]
└── MergeExecutor { output: [ team_to_accounts_mv.team_id, team_to_accounts_mv.account_id ], stream key: [ team_to_accounts_mv.team_id, team_to_accounts_mv.account_id ] }
Fragment 52689 (Actor 738234,738233)
StreamTableScan { table: team_to_accounts_mv, columns: [team_id, account_id] }
├── output: [ team_to_accounts_mv.team_id, team_to_accounts_mv.account_id ]
├── stream key: [ team_to_accounts_mv.team_id, team_to_accounts_mv.account_id ]
├── Upstream { output: [ team_id, account_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, account_id ], stream key: [] }
Fragment 52690 (Actor 738194,738195)
StreamProject { exprs: [team_to_portfolios_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id, 2:Int32] }
├── output: [ team_to_portfolios_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id, 2:Int32 ]
├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id ]
└── MergeExecutor
├── output: [ team_to_portfolios_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, team_to_portfolios_mv.portfolio_id ]
└── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id ]
Fragment 52691 (Actor 738197,738196)
StreamSyncLogStore
├── output: [ team_to_portfolios_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, team_to_portfolios_mv.portfolio_id ]
├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id ]
└── StreamHashJoin { type: Inner, predicate: tasks_dm.resource_portfolio_id = team_to_portfolios_mv.portfolio_id }
├── output: [ team_to_portfolios_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_portfolio_id, team_to_portfolios_mv.portfolio_id ]
├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id, team_to_portfolios_mv.team_id ]
├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at ], stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id ] }
└── MergeExecutor
├── output: [ team_to_portfolios_mv.team_id, team_to_portfolios_mv.portfolio_id ]
└── stream key: [ team_to_portfolios_mv.portfolio_id, team_to_portfolios_mv.team_id ]
Fragment 52692 (Actor 738199,738198)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_portfolio_id] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.resource_portfolio_id, tasks_dm.task_id ]
└── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_portfolio_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }
Fragment 52693 (Actor 738236,738235)
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 52694 (Actor 738204,738203)
StreamLocalityProvider { locality_columns: [team_to_portfolios_mv.portfolio_id] }
├── output: [ team_to_portfolios_mv.team_id, team_to_portfolios_mv.portfolio_id ]
├── stream key: [ team_to_portfolios_mv.portfolio_id, team_to_portfolios_mv.team_id ]
└── MergeExecutor { output: [ team_to_portfolios_mv.team_id, team_to_portfolios_mv.portfolio_id ], stream key: [ team_to_portfolios_mv.team_id, team_to_portfolios_mv.portfolio_id ] }
Fragment 52695 (Actor 738237,738238)
StreamTableScan { table: team_to_portfolios_mv, columns: [team_id, portfolio_id] }
├── output: [ team_to_portfolios_mv.team_id, team_to_portfolios_mv.portfolio_id ]
├── stream key: [ team_to_portfolios_mv.team_id, team_to_portfolios_mv.portfolio_id ]
├── Upstream { output: [ team_id, portfolio_id ], stream key: [] }
└── BatchPlanNode { output: [ team_id, portfolio_id ], stream key: [] }
Fragment 52696 (Actor 738207,738208)
StreamProject { exprs: [team_to_parties_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_id, 3:Int32] }
├── output: [ team_to_parties_mv.team_id, tasks_dm.task_id, tasks_dm.updated_at, tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_id, 3:Int32 ]
├── stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_id ]
└── MergeExecutor
├── 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.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_id ]
Fragment 52697 (Actor 738209,738210)
StreamSyncLogStore
├── 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.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_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.resource_party_id, tasks_dm.task_id, team_to_parties_mv.team_id ]
├── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at ], stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id ] }
└── MergeExecutor { output: [ team_to_parties_mv.team_id, team_to_parties_mv.party_id ], stream key: [ team_to_parties_mv.party_id, team_to_parties_mv.team_id ] }
Fragment 52698 (Actor 738211,738212)
StreamLocalityProvider { locality_columns: [tasks_dm.resource_party_id] }
├── output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at ]
├── stream key: [ tasks_dm.resource_party_id, tasks_dm.task_id ]
└── MergeExecutor { output: [ tasks_dm.task_id, tasks_dm.resource_party_id, tasks_dm.updated_at ], stream key: [ tasks_dm.task_id ] }
Fragment 52699 (Actor 738239,738240)
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 52700 (Actor 738219,738220)
StreamLocalityProvider { locality_columns: [team_to_parties_mv.party_id] }
├── output: [ team_to_parties_mv.team_id, team_to_parties_mv.party_id ]
├── stream key: [ team_to_parties_mv.party_id, team_to_parties_mv.team_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 52701 (Actor 738221,738222)
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: [] }