Job is idle — throughput ~0; structure shown.
Fragment 51736 (Actor 736613,736612)
StreamMaterialize { columns: [account_id, active_count, high_priority_count], stream_key: [account_id], pk_columns: [account_id], pk_conflict: NoCheck } { output: [ tasks_dm.resource_account_id, $expr1, $expr2 ], stream key: [ tasks_dm.resource_account_id ] }
└── StreamProject { exprs: [tasks_dm.resource_account_id, count filter(In(tasks_dm.status, 'TO_DO':Varchar, 'IN_PROGRESS':Varchar))::Int32 as $expr1, count filter(In(tasks_dm.status, 'TO_DO':Varchar, 'IN_PROGRESS':Varchar) AND (tasks_dm.priority = 'HIGH':Varchar))::Int32 as $expr2] }
├── output: [ tasks_dm.resource_account_id, $expr1, $expr2 ]
├── stream key: [ tasks_dm.resource_account_id ]
└── StreamHashAgg { group_key: [tasks_dm.resource_account_id], aggs: [count filter(In(tasks_dm.status, 'TO_DO':Varchar, 'IN_PROGRESS':Varchar)), count filter(In(tasks_dm.status, 'TO_DO':Varchar, 'IN_PROGRESS':Varchar) AND (tasks_dm.priority = 'HIGH':Varchar)), count] }
├── output: [ tasks_dm.resource_account_id, count filter(In(tasks_dm.status, 'TO_DO':Varchar, 'IN_PROGRESS':Varchar)), count filter(In(tasks_dm.status, 'TO_DO':Varchar, 'IN_PROGRESS':Varchar) AND (tasks_dm.priority = 'HIGH':Varchar)), count ]
├── stream key: [ tasks_dm.resource_account_id ]
└── StreamLocalityProvider { locality_columns: [tasks_dm.resource_account_id] } { output: [ tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, tasks_dm.task_id ], stream key: [ tasks_dm.resource_account_id, tasks_dm.task_id ] }
└── MergeExecutor { output: [ tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, tasks_dm.task_id ], stream key: [ tasks_dm.task_id ] }
Fragment 51737 (Actor 736614,736615)
StreamProject { exprs: [tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, tasks_dm.task_id] } { output: [ tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, tasks_dm.task_id ], stream key: [ tasks_dm.task_id ] }
└── StreamFilter { predicate: Not(IsNull(tasks_dm.resource_account_id)) AND IsNull(tasks_dm.disabled_at) } { output: [ tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, tasks_dm.task_id, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
└── StreamTableScan { table: tasks_dm, columns: [status, priority, resource_account_id, task_id, disabled_at] } { output: [ tasks_dm.status, tasks_dm.priority, tasks_dm.resource_account_id, tasks_dm.task_id, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
├── Upstream { output: [ status, priority, resource_account_id, task_id, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ status, priority, resource_account_id, task_id, disabled_at ], stream key: [] }