Job is idle — throughput ~0; structure shown.
Fragment 44390 (Actor 741801,741800)
StreamMaterialize { columns: [client_id, active_count, high_priority_count], stream_key: [client_id], pk_columns: [client_id], pk_conflict: NoCheck } { output: [ tasks_dm.resource_client_id, $expr1, $expr2 ], stream key: [ tasks_dm.resource_client_id ] }
└── StreamProject { exprs: [tasks_dm.resource_client_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_client_id, $expr1, $expr2 ]
├── stream key: [ tasks_dm.resource_client_id ]
└── StreamHashAgg { group_key: [tasks_dm.resource_client_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_client_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_client_id ]
└── StreamLocalityProvider { locality_columns: [tasks_dm.resource_client_id] } { output: [ tasks_dm.status, tasks_dm.priority, tasks_dm.resource_client_id, tasks_dm.task_id ], stream key: [ tasks_dm.resource_client_id, tasks_dm.task_id ] }
└── MergeExecutor { output: [ tasks_dm.status, tasks_dm.priority, tasks_dm.resource_client_id, tasks_dm.task_id ], stream key: [ tasks_dm.task_id ] }
Fragment 44391 (Actor 741827,741828)
StreamProject { exprs: [tasks_dm.status, tasks_dm.priority, tasks_dm.resource_client_id, tasks_dm.task_id] } { output: [ tasks_dm.status, tasks_dm.priority, tasks_dm.resource_client_id, tasks_dm.task_id ], stream key: [ tasks_dm.task_id ] }
└── StreamFilter { predicate: Not(IsNull(tasks_dm.resource_client_id)) AND IsNull(tasks_dm.disabled_at) } { output: [ tasks_dm.status, tasks_dm.priority, tasks_dm.resource_client_id, tasks_dm.task_id, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
└── StreamTableScan { table: tasks_dm, columns: [status, priority, resource_client_id, task_id, disabled_at] } { output: [ tasks_dm.status, tasks_dm.priority, tasks_dm.resource_client_id, tasks_dm.task_id, tasks_dm.disabled_at ], stream key: [ tasks_dm.task_id ] }
├── Upstream { output: [ status, priority, resource_client_id, task_id, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ status, priority, resource_client_id, task_id, disabled_at ], stream key: [] }