Job is idle — throughput ~0; structure shown.
Fragment 51604 (Actor 736019,736018)
StreamSink { type: upsert, columns: [account_id, base_currency_code, is_restricted, product_type, accounts_dm.product_type_id(hidden), product_types_dm.product_type_id(hidden)], downstream_pk: [accounts_dm.account_id] }
├── output: [ accounts_dm.account_id, accounts_dm.base_currency_code, accounts_dm.is_restricted, product_types_dm.type, accounts_dm.product_type_id, product_types_dm.product_type_id ]
├── stream key: [ accounts_dm.product_type_id, accounts_dm.account_id ]
└── MergeExecutor
├── output: [ accounts_dm.account_id, accounts_dm.base_currency_code, accounts_dm.is_restricted, product_types_dm.type, accounts_dm.product_type_id, product_types_dm.product_type_id ]
└── stream key: [ accounts_dm.product_type_id, accounts_dm.account_id ]
Fragment 51605 (Actor 736083,736084)
StreamSyncLogStore
├── output: [ accounts_dm.account_id, accounts_dm.base_currency_code, accounts_dm.is_restricted, product_types_dm.type, accounts_dm.product_type_id, product_types_dm.product_type_id ]
├── stream key: [ accounts_dm.product_type_id, accounts_dm.account_id ]
└── StreamHashJoin { type: LeftOuter, predicate: accounts_dm.product_type_id = product_types_dm.product_type_id }
├── output: [ accounts_dm.account_id, accounts_dm.base_currency_code, accounts_dm.is_restricted, product_types_dm.type, accounts_dm.product_type_id, product_types_dm.product_type_id ]
├── stream key: [ accounts_dm.product_type_id, accounts_dm.account_id ]
├── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.is_restricted ], stream key: [ accounts_dm.product_type_id, accounts_dm.account_id ] }
└── MergeExecutor { output: [ product_types_dm.product_type_id, product_types_dm.type ], stream key: [ product_types_dm.product_type_id ] }
Fragment 51606 (Actor 736086,736085)
StreamLocalityProvider { locality_columns: [accounts_dm.product_type_id] }
├── output: [ accounts_dm.account_id, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.is_restricted ]
├── stream key: [ accounts_dm.product_type_id, accounts_dm.account_id ]
└── MergeExecutor { output: [ accounts_dm.account_id, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.is_restricted ], stream key: [ accounts_dm.account_id ] }
Fragment 51607 (Actor 736088,736087)
StreamProject { exprs: [accounts_dm.account_id, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.is_restricted] }
├── output: [ accounts_dm.account_id, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.is_restricted ]
├── stream key: [ accounts_dm.account_id ]
└── StreamFilter { predicate: IsNull(accounts_dm.closing_date) AND IsNull(accounts_dm.disabled_at) }
├── output: [ accounts_dm.account_id, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.is_restricted, accounts_dm.closing_date, accounts_dm.disabled_at ]
├── stream key: [ accounts_dm.account_id ]
└── StreamTableScan { table: accounts_dm, columns: [account_id, product_type_id, base_currency_code, is_restricted, closing_date, disabled_at] }
├── output: [ accounts_dm.account_id, accounts_dm.product_type_id, accounts_dm.base_currency_code, accounts_dm.is_restricted, accounts_dm.closing_date, accounts_dm.disabled_at ]
├── stream key: [ accounts_dm.account_id ]
├── Upstream { output: [ account_id, product_type_id, base_currency_code, is_restricted, closing_date, disabled_at ], stream key: [] }
└── BatchPlanNode { output: [ account_id, product_type_id, base_currency_code, is_restricted, closing_date, disabled_at ], stream key: [] }
Fragment 51608 (Actor 736075,736076)
StreamTableScan { table: product_types_dm, columns: [product_type_id, type] } { output: [ product_types_dm.product_type_id, product_types_dm.type ], stream key: [ product_types_dm.product_type_id ] }
├── Upstream { output: [ product_type_id, type ], stream key: [] }
└── BatchPlanNode { output: [ product_type_id, type ], stream key: [] }