Job is idle — throughput ~0; structure shown.
Fragment 37109 (Actor 742611,742612)
StreamMaterialize { columns: [id, client_order_id, portfolio_id, security_account_id, initiator, asset_id, side, workflow_state, execution_state, etag, customer_relationship_id, instruction, execution_spec, cost_estimate, filled_quantity, created_at, deleted_at, termination_reason, executed_at, termination_detail], stream_key: [id, created_at], pk_columns: [id], pk_conflict: Overwrite, watermark_columns: [created_at] }
├── output:
│ ┌── oms.orders.id
│ ├── oms.orders.client_order_id
│ ├── oms.orders.portfolio_id
│ ├── oms.orders.security_account_id
│ ├── oms.orders.initiator
│ ├── oms.orders.asset_id
│ ├── oms.orders.side
│ ├── oms.orders.workflow_state
│ ├── oms.orders.execution_state
│ ├── oms.orders.etag
│ ├── oms.orders.customer_relationship_id
│ ├── oms.orders.instruction
│ ├── oms.orders.execution_spec
│ ├── oms.orders.cost_estimate
│ ├── oms.orders.filled_quantity
│ ├── oms.orders.created_at
│ ├── oms.orders.deleted_at
│ ├── oms.orders.termination_reason
│ ├── oms.orders.executed_at
│ └── oms.orders.termination_detail
├── stream key: [ oms.orders.id, oms.orders.created_at ]
└── StreamWatermarkFilter [upsert] { watermark_descs: [Desc { column: oms.orders.created_at, expr: SubtractWithTimeZone(oms.orders.created_at, '5 years':Interval, 'UTC':Varchar) }], output_watermarks: [[oms.orders.created_at]] }
├── output:
│ ┌── oms.orders.id
│ ├── oms.orders.client_order_id
│ ├── oms.orders.portfolio_id
│ ├── oms.orders.security_account_id
│ ├── oms.orders.initiator
│ ├── oms.orders.asset_id
│ ├── oms.orders.side
│ ├── oms.orders.workflow_state
│ ├── oms.orders.execution_state
│ ├── oms.orders.etag
│ ├── oms.orders.customer_relationship_id
│ ├── oms.orders.instruction
│ ├── oms.orders.execution_spec
│ ├── oms.orders.cost_estimate
│ ├── oms.orders.filled_quantity
│ ├── oms.orders.created_at
│ ├── oms.orders.deleted_at
│ ├── oms.orders.termination_reason
│ ├── oms.orders.executed_at
│ └── oms.orders.termination_detail
├── stream key: []
└── StreamUnion { all: true }
├── output:
│ ┌── oms.orders.id
│ ├── oms.orders.client_order_id
│ ├── oms.orders.portfolio_id
│ ├── oms.orders.security_account_id
│ ├── oms.orders.initiator
│ ├── oms.orders.asset_id
│ ├── oms.orders.side
│ ├── oms.orders.workflow_state
│ ├── oms.orders.execution_state
│ ├── oms.orders.etag
│ ├── oms.orders.customer_relationship_id
│ ├── oms.orders.instruction
│ ├── oms.orders.execution_spec
│ ├── oms.orders.cost_estimate
│ ├── oms.orders.filled_quantity
│ ├── oms.orders.created_at
│ ├── oms.orders.deleted_at
│ ├── oms.orders.termination_reason
│ ├── oms.orders.executed_at
│ └── oms.orders.termination_detail
├── stream key: []
├── MergeExecutor
│ ├── output:
│ │ ┌── oms.orders.id
│ │ ├── oms.orders.client_order_id
│ │ ├── oms.orders.portfolio_id
│ │ ├── oms.orders.security_account_id
│ │ ├── oms.orders.initiator
│ │ ├── oms.orders.asset_id
│ │ ├── oms.orders.side
│ │ ├── oms.orders.workflow_state
│ │ ├── oms.orders.execution_state
│ │ ├── oms.orders.etag
│ │ ├── oms.orders.customer_relationship_id
│ │ ├── oms.orders.instruction
│ │ ├── oms.orders.execution_spec
│ │ ├── oms.orders.cost_estimate
│ │ ├── oms.orders.filled_quantity
│ │ ├── oms.orders.created_at
│ │ ├── oms.orders.deleted_at
│ │ ├── oms.orders.termination_reason
│ │ ├── oms.orders.executed_at
│ │ └── oms.orders.termination_detail
│ └── stream key: [ oms.orders.id ]
├── MergeExecutor { output: [ id, client_order_id, portfolio_id, security_account_id, initiator, asset_id, side, workflow_state, execution_state, etag, customer_relationship_id, instruction, execution_spec, cost_estimate, filled_quantity, created_at, deleted_at, termination_reason, executed_at, termination_detail ], stream key: [] }
└── StreamUpstreamSinkUnion { output: [ id, client_order_id, portfolio_id, security_account_id, initiator, asset_id, side, workflow_state, execution_state, etag, customer_relationship_id, instruction, execution_spec, cost_estimate, filled_quantity, created_at, deleted_at, termination_reason, executed_at, termination_detail ], stream key: [] }
Fragment 37110 (Actor 742613)
StreamCdcTableScan { table: oms.orders, columns: [id, client_order_id, portfolio_id, security_account_id, initiator, asset_id, side, workflow_state, execution_state, etag, customer_relationship_id, instruction, execution_spec, cost_estimate, filled_quantity, created_at, deleted_at, termination_reason, executed_at, termination_detail] }
├── output:
│ ┌── oms.orders.id
│ ├── oms.orders.client_order_id
│ ├── oms.orders.portfolio_id
│ ├── oms.orders.security_account_id
│ ├── oms.orders.initiator
│ ├── oms.orders.asset_id
│ ├── oms.orders.side
│ ├── oms.orders.workflow_state
│ ├── oms.orders.execution_state
│ ├── oms.orders.etag
│ ├── oms.orders.customer_relationship_id
│ ├── oms.orders.instruction
│ ├── oms.orders.execution_spec
│ ├── oms.orders.cost_estimate
│ ├── oms.orders.filled_quantity
│ ├── oms.orders.created_at
│ ├── oms.orders.deleted_at
│ ├── oms.orders.termination_reason
│ ├── oms.orders.executed_at
│ └── oms.orders.termination_detail
├── stream key: [ oms.orders.id ]
└── MergeExecutor { output: [ payload, _rw_offset, _rw_table_name ], stream key: [] }
Fragment 37111 (Actor 741824)
StreamCdcFilter { output: [ payload, _rw_offset, _rw_table_name ], stream key: [] }
└── Upstream { output: [ payload, _rw_offset, _rw_table_name ], stream key: [] }
Fragment 37112 (Actor 742615,742614)
StreamDml { columns: [id, client_order_id, portfolio_id, security_account_id, initiator, asset_id, side, workflow_state, execution_state, etag, customer_relationship_id, instruction, execution_spec, cost_estimate, filled_quantity, created_at, deleted_at, termination_reason, executed_at, termination_detail] }
├── output: [ id, client_order_id, portfolio_id, security_account_id, initiator, asset_id, side, workflow_state, execution_state, etag, customer_relationship_id, instruction, execution_spec, cost_estimate, filled_quantity, created_at, deleted_at, termination_reason, executed_at, termination_detail ]
├── stream key: []
└── StreamSource { output: [ id, client_order_id, portfolio_id, security_account_id, initiator, asset_id, side, workflow_state, execution_state, etag, customer_relationship_id, instruction, execution_spec, cost_estimate, filled_quantity, created_at, deleted_at, termination_reason, executed_at, termination_detail ], stream key: [] }