Fragment 25027 (Actor 119902,119903)
StreamMaterialize { columns: [account_id, asset_id, currency_code, dim_value_date, purchased_quantity, average_cost_per_unit, market_value], stream_key: [account_id, asset_id], pk_columns: [account_id, asset_id], pk_conflict: NoCheck }
├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.currency_code, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.market_value ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id ]
└── StreamProject { exprs: [holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.currency_code, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.market_value] }
├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.currency_code, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.market_value ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id ]
└── StreamGroupTopN { order: [holding_values_raw_ft.dim_value_date DESC], limit: 1, offset: 0, group_key: [holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id] }
├── output:
│ ┌── holding_values_raw_ft.account_id
│ ├── holding_values_raw_ft.asset_id
│ ├── holding_values_raw_ft.dim_value_date
│ ├── holding_values_raw_ft.currency_code
│ ├── holding_values_raw_ft.market_value
│ ├── holding_values_raw_ft.average_cost_per_unit
│ ├── holding_values_raw_ft.purchased_quantity
│ └── holding_values_raw_ft.type
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id ]
└── MergeExecutor
├── output:
│ ┌── holding_values_raw_ft.account_id
│ ├── holding_values_raw_ft.asset_id
│ ├── holding_values_raw_ft.dim_value_date
│ ├── holding_values_raw_ft.currency_code
│ ├── holding_values_raw_ft.market_value
│ ├── holding_values_raw_ft.average_cost_per_unit
│ ├── holding_values_raw_ft.purchased_quantity
│ └── holding_values_raw_ft.type
└── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type ]
Fragment 25028 (Actor 119906,119905)
StreamProject { exprs: [holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.type] }
├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.type ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type ]
└── StreamDynamicFilter { predicate: ($expr1 <= now), output: [holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, $expr1, holding_values_raw_ft.type], cleaned_by_watermark: true }
├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, $expr1, holding_values_raw_ft.type ]
├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type ]
├── StreamProject { exprs: [holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, AtTimeZone(holding_values_raw_ft.dim_value_date::Timestamp, 'UTC':Varchar) as $expr1, holding_values_raw_ft.type], output_watermarks: [[holding_values_raw_ft.dim_value_date]] }
│ ├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, $expr1, holding_values_raw_ft.type ]
│ ├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type ]
│ └── StreamFilter { predicate: (holding_values_raw_ft.type = 'ASSET':Varchar) AND IsNull(holding_values_raw_ft.disabled_at) AND (Coalesce(holding_values_raw_ft.m_is_stub, false:Boolean) = false:Boolean) }
│ ├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.type, holding_values_raw_ft.m_is_stub, holding_values_raw_ft.disabled_at ]
│ ├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type ]
│ └── StreamTableScan { table: holding_values_raw_ft, columns: [account_id, asset_id, dim_value_date, currency_code, market_value, average_cost_per_unit, purchased_quantity, type, m_is_stub, disabled_at] }
│ ├── output: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.currency_code, holding_values_raw_ft.market_value, holding_values_raw_ft.average_cost_per_unit, holding_values_raw_ft.purchased_quantity, holding_values_raw_ft.type, holding_values_raw_ft.m_is_stub, holding_values_raw_ft.disabled_at ]
│ ├── stream key: [ holding_values_raw_ft.account_id, holding_values_raw_ft.asset_id, holding_values_raw_ft.dim_value_date, holding_values_raw_ft.type ]
│ ├── Upstream { output: [ account_id, asset_id, dim_value_date, currency_code, market_value, average_cost_per_unit, purchased_quantity, type, m_is_stub, disabled_at ], stream key: [] }
│ └── BatchPlanNode { output: [ account_id, asset_id, dim_value_date, currency_code, market_value, average_cost_per_unit, purchased_quantity, type, m_is_stub, disabled_at ], stream key: [] }
└── MergeExecutor { output: [ now ], stream key: [] }
Fragment 25029 (Actor 119904)
StreamNow { output: [ now ], stream key: [] }