Fragment 25087 (Actor 120530,120529)
StreamMaterialize { columns: [id, security_account_id, portfolio_id, created_at, order_side, order_status, quantity, filled_quantity, average_price, gross_executed_amount, total_fee, estimated_gross_amount, estimated_net_amount, estimated_fee_amount, currency_code, asset_id, asset_name_en, asset_name_ar, asset_type, asset_currency_code, asset_ticker, asset_isin, orders.asset_id(hidden)], stream_key: [id, created_at, orders.asset_id], pk_columns: [id, created_at, orders.asset_id], pk_conflict: NoCheck } { output: [ orders.id, orders.security_account_id, orders.portfolio_id, orders.created_at, $expr2, $expr3, $expr4, $expr5, $expr6, $expr7, $expr8, $expr9, $expr10, $expr11, $expr12, assets_dm_next.id, assets_dm_next.name_en, assets_dm_next.name_ar, assets_dm_next.type, assets_dm_next.issue_currency_code, assets_dm_next.ticker, assets_dm_next.isin, orders.asset_id ], stream key: [ orders.id, orders.created_at, orders.asset_id ] }
└── StreamProject
└─exprs:
┌─orders.id
├─orders.security_account_id
├─orders.portfolio_id
├─orders.created_at
├─Case((orders.side = 'TRADE_SIDE_BUY':Varchar), 'BUY':Varchar, (orders.side = 'TRADE_SIDE_SELL':Varchar), 'SELL':Varchar, (orders.side = 'TRADE_SIDE_SUBSCRIBE':Varchar), 'SUBSCRIBE':Varchar, (orders.side = 'TRADE_SIDE_REDEEM':Varchar), 'REDEEM':Varchar, (orders.side = 'TRADE_SIDE_DRIP':Varchar), 'DRIP':Varchar, orders.side) as $expr2
├─Case((orders.workflow_state = 'WORKFLOW_STATE_DRAFT':Varchar), 'DRAFT':Varchar, (orders.workflow_state = 'WORKFLOW_STATE_PENDING_COMPLIANCE':Varchar), 'PENDING_COMPLIANCE':Varchar, (orders.workflow_state = 'WORKFLOW_STATE_PENDING_APPROVAL':Varchar), 'PENDING_APPROVAL':Varchar, (orders.workflow_state = 'WORKFLOW_STATE_REQUIRES_MODIFICATION':Varchar), 'REQUIRES_MODIFICATION':Varchar, (orders.workflow_state = 'WORKFLOW_STATE_TERMINAL':Varchar), Case((orders.execution_state = 'EXECUTION_STATE_FILLED':Varchar), 'FILLED':Varchar, (orders.execution_state = 'EXECUTION_STATE_PARTIALLY_FILLED':Varchar), 'CANCELED':Varchar, (orders.execution_state = 'EXECUTION_STATE_CANCELLED':Varchar), 'CANCELED':Varchar, (orders.execution_state = 'EXECUTION_STATE_REJECTED':Varchar), 'REJECTED':Varchar, (orders.execution_state = 'EXECUTION_STATE_EXPIRED':Varchar), 'EXPIRED':Varchar, null:Varchar), (orders.workflow_state = 'WORKFLOW_STATE_ACTIVE':Varchar), Case((orders.execution_state = 'EXECUTION_STATE_NEW':Varchar), 'NEW':Varchar, (orders.execution_state = 'EXECUTION_STATE_ACKNOWLEDGED':Varchar), 'NEW':Varchar, (orders.execution_state = 'EXECUTION_STATE_PARTIALLY_FILLED':Varchar), 'PARTIALLY_FILLED':Varchar, (orders.execution_state = 'EXECUTION_STATE_PENDING_CANCEL':Varchar), 'PENDING_CANCEL':Varchar, (orders.execution_state = 'EXECUTION_STATE_FILLED':Varchar), 'FILLED':Varchar, 'NEW':Varchar), null:Varchar) as $expr3
├─Coalesce(JsonbAccessStr(JsonbAccess(JsonbAccess(orders.execution_spec, 'order_book':Varchar), 'shares':Varchar), 'value':Varchar)::Decimal, JsonbAccessStr(JsonbAccess(JsonbAccess(JsonbAccess(orders.execution_spec, 'order_book':Varchar), 'money':Varchar), 'amount':Varchar), 'value':Varchar)::Decimal) as $expr4
├─Coalesce(sdk_order_execution_aggregates_mv.filled_quantity, 0:Decimal) as $expr5
├─Coalesce(sdk_order_execution_aggregates_mv.average_price, 0:Decimal) as $expr6
├─Coalesce(sdk_order_execution_aggregates_mv.total_amount, 0:Decimal) as $expr7
├─Coalesce(sdk_order_execution_aggregates_mv.total_fees, 0:Decimal) as $expr8
├─JsonbAccessStr(JsonbAccess(JsonbAccess(orders.cost_estimate, 'estimatedGross':Varchar), 'amount':Varchar), 'value':Varchar)::Decimal as $expr9
├─JsonbAccessStr(JsonbAccess(JsonbAccess(orders.cost_estimate, 'estimatedNet':Varchar), 'amount':Varchar), 'value':Varchar)::Decimal as $expr10
├─JsonbAccessStr(JsonbAccess(JsonbAccess(orders.cost_estimate, 'estimatedFees':Varchar), 'amount':Varchar), 'value':Varchar)::Decimal as $expr11
├─Coalesce(sdk_order_execution_aggregates_mv.execution_currency_code, JsonbAccessStr(JsonbAccess(JsonbAccess(orders.cost_estimate, 'estimatedGross':Varchar), 'currencyCode':Varchar), 'value':Varchar), JsonbAccessStr(JsonbAccess(JsonbAccess(JsonbAccess(orders.execution_spec, 'order_book':Varchar), 'limit_price':Varchar), 'currency_code':Varchar), 'value':Varchar), JsonbAccessStr(JsonbAccess(JsonbAccess(JsonbAccess(orders.execution_spec, 'order_book':Varchar), 'money':Varchar), 'currency_code':Varchar), 'value':Varchar), assets_dm_next.issue_currency_code) as $expr12
├─assets_dm_next.id
├─assets_dm_next.name_en
├─assets_dm_next.name_ar
├─assets_dm_next.type
├─assets_dm_next.issue_currency_code
├─assets_dm_next.ticker
├─assets_dm_next.isin
└─orders.asset_id
├── output: [ orders.id, orders.security_account_id, orders.portfolio_id, orders.created_at, $expr2, $expr3, $expr4, $expr5, $expr6, $expr7, $expr8, $expr9, $expr10, $expr11, $expr12, assets_dm_next.id, assets_dm_next.name_en, assets_dm_next.name_ar, assets_dm_next.type, assets_dm_next.issue_currency_code, assets_dm_next.ticker, assets_dm_next.isin, orders.asset_id ]
├── stream key: [ orders.id, orders.created_at, orders.asset_id ]
└── StreamDynamicFilter { predicate: (orders.created_at >= $expr1), output: [orders.id, orders.portfolio_id, orders.security_account_id, orders.side, orders.workflow_state, orders.execution_state, orders.execution_spec, orders.cost_estimate, orders.created_at, sdk_order_execution_aggregates_mv.filled_quantity, sdk_order_execution_aggregates_mv.total_amount, sdk_order_execution_aggregates_mv.total_fees, sdk_order_execution_aggregates_mv.average_price, sdk_order_execution_aggregates_mv.execution_currency_code, assets_dm_next.id, assets_dm_next.name_en, assets_dm_next.name_ar, assets_dm_next.type, assets_dm_next.issue_currency_code, assets_dm_next.ticker, assets_dm_next.isin, orders.asset_id] } { output: [ orders.id, orders.portfolio_id, orders.security_account_id, orders.side, orders.workflow_state, orders.execution_state, orders.execution_spec, orders.cost_estimate, orders.created_at, sdk_order_execution_aggregates_mv.filled_quantity, sdk_order_execution_aggregates_mv.total_amount, sdk_order_execution_aggregates_mv.total_fees, sdk_order_execution_aggregates_mv.average_price, sdk_order_execution_aggregates_mv.execution_currency_code, assets_dm_next.id, assets_dm_next.name_en, assets_dm_next.name_ar, assets_dm_next.type, assets_dm_next.issue_currency_code, assets_dm_next.ticker, assets_dm_next.isin, orders.asset_id ], stream key: [ orders.id, orders.created_at, orders.asset_id ] }
├── StreamTemporalJoin { type: LeftOuter, append_only: false, predicate: orders.asset_id = assets_dm_next.id, nested_loop: false } { output: [ orders.id, orders.portfolio_id, orders.security_account_id, orders.side, orders.workflow_state, orders.execution_state, orders.execution_spec, orders.cost_estimate, orders.created_at, sdk_order_execution_aggregates_mv.filled_quantity, sdk_order_execution_aggregates_mv.total_amount, sdk_order_execution_aggregates_mv.total_fees, sdk_order_execution_aggregates_mv.average_price, sdk_order_execution_aggregates_mv.execution_currency_code, assets_dm_next.id, assets_dm_next.name_en, assets_dm_next.name_ar, assets_dm_next.type, assets_dm_next.issue_currency_code, assets_dm_next.ticker, assets_dm_next.isin, orders.asset_id ], stream key: [ orders.id, orders.created_at, orders.asset_id ] }
│ ├── MergeExecutor { output: [ orders.id, orders.portfolio_id, orders.security_account_id, orders.asset_id, orders.side, orders.workflow_state, orders.execution_state, orders.execution_spec, orders.cost_estimate, orders.filled_quantity, orders.created_at, sdk_order_execution_aggregates_mv.filled_quantity, sdk_order_execution_aggregates_mv.total_amount, sdk_order_execution_aggregates_mv.total_fees, sdk_order_execution_aggregates_mv.average_price, sdk_order_execution_aggregates_mv.execution_currency_code, sdk_order_execution_aggregates_mv.order_id ], stream key: [ orders.id, orders.created_at ] }
│ └── MergeExecutor { output: [ assets_dm_next.id, assets_dm_next.name_en, assets_dm_next.name_ar, assets_dm_next.type, assets_dm_next.issue_currency_code, assets_dm_next.ticker, assets_dm_next.isin ], stream key: [ assets_dm_next.id ] }
└── MergeExecutor { output: [ $expr1 ], stream key: [] }
Fragment 25088 (Actor 120556,120557)
StreamFilter { predicate: ((Not((orders.workflow_state = 'WORKFLOW_STATE_TERMINAL':Varchar)) OR Not((Coalesce(orders.filled_quantity, 0:Decimal) > 0:Decimal))) OR (Coalesce(sdk_order_execution_aggregates_mv.filled_quantity, 0:Decimal) = Coalesce(orders.filled_quantity, 0:Decimal))) } { output: [ orders.id, orders.portfolio_id, orders.security_account_id, orders.asset_id, orders.side, orders.workflow_state, orders.execution_state, orders.execution_spec, orders.cost_estimate, orders.filled_quantity, orders.created_at, sdk_order_execution_aggregates_mv.filled_quantity, sdk_order_execution_aggregates_mv.total_amount, sdk_order_execution_aggregates_mv.total_fees, sdk_order_execution_aggregates_mv.average_price, sdk_order_execution_aggregates_mv.execution_currency_code, sdk_order_execution_aggregates_mv.order_id ], stream key: [ orders.id, orders.created_at ] }
└── MergeExecutor { output: [ orders.id, orders.portfolio_id, orders.security_account_id, orders.asset_id, orders.side, orders.workflow_state, orders.execution_state, orders.execution_spec, orders.cost_estimate, orders.filled_quantity, orders.created_at, sdk_order_execution_aggregates_mv.filled_quantity, sdk_order_execution_aggregates_mv.total_amount, sdk_order_execution_aggregates_mv.total_fees, sdk_order_execution_aggregates_mv.average_price, sdk_order_execution_aggregates_mv.execution_currency_code, sdk_order_execution_aggregates_mv.order_id ], stream key: [ orders.id, orders.created_at ] }
Fragment 25089 (Actor 120558,120559)
StreamSyncLogStore { output: [ orders.id, orders.portfolio_id, orders.security_account_id, orders.asset_id, orders.side, orders.workflow_state, orders.execution_state, orders.execution_spec, orders.cost_estimate, orders.filled_quantity, orders.created_at, sdk_order_execution_aggregates_mv.filled_quantity, sdk_order_execution_aggregates_mv.total_amount, sdk_order_execution_aggregates_mv.total_fees, sdk_order_execution_aggregates_mv.average_price, sdk_order_execution_aggregates_mv.execution_currency_code, sdk_order_execution_aggregates_mv.order_id ], stream key: [ orders.id, orders.created_at ] }
└── StreamHashJoin { type: LeftOuter, predicate: orders.id = sdk_order_execution_aggregates_mv.order_id } { output: [ orders.id, orders.portfolio_id, orders.security_account_id, orders.asset_id, orders.side, orders.workflow_state, orders.execution_state, orders.execution_spec, orders.cost_estimate, orders.filled_quantity, orders.created_at, sdk_order_execution_aggregates_mv.filled_quantity, sdk_order_execution_aggregates_mv.total_amount, sdk_order_execution_aggregates_mv.total_fees, sdk_order_execution_aggregates_mv.average_price, sdk_order_execution_aggregates_mv.execution_currency_code, sdk_order_execution_aggregates_mv.order_id ], stream key: [ orders.id, orders.created_at ] }
├── MergeExecutor { output: [ orders.id, orders.portfolio_id, orders.security_account_id, orders.asset_id, orders.side, orders.workflow_state, orders.execution_state, orders.execution_spec, orders.cost_estimate, orders.filled_quantity, orders.created_at ], stream key: [ orders.id, orders.created_at ] }
└── MergeExecutor { output: [ sdk_order_execution_aggregates_mv.order_id, sdk_order_execution_aggregates_mv.filled_quantity, sdk_order_execution_aggregates_mv.total_amount, sdk_order_execution_aggregates_mv.total_fees, sdk_order_execution_aggregates_mv.average_price, sdk_order_execution_aggregates_mv.execution_currency_code ], stream key: [ sdk_order_execution_aggregates_mv.order_id ] }
Fragment 25090 (Actor 120560,120561)
StreamTableScan { table: orders, columns: [id, portfolio_id, security_account_id, asset_id, side, workflow_state, execution_state, execution_spec, cost_estimate, filled_quantity, created_at] } { output: [ orders.id, orders.portfolio_id, orders.security_account_id, orders.asset_id, orders.side, orders.workflow_state, orders.execution_state, orders.execution_spec, orders.cost_estimate, orders.filled_quantity, orders.created_at ], stream key: [ orders.id, orders.created_at ] }
├── Upstream { output: [ id, portfolio_id, security_account_id, asset_id, side, workflow_state, execution_state, execution_spec, cost_estimate, filled_quantity, created_at ], stream key: [] }
└── BatchPlanNode { output: [ id, portfolio_id, security_account_id, asset_id, side, workflow_state, execution_state, execution_spec, cost_estimate, filled_quantity, created_at ], stream key: [] }
Fragment 25091 (Actor 120563,120562)
StreamTableScan { table: sdk_order_execution_aggregates_mv, columns: [order_id, filled_quantity, total_amount, total_fees, average_price, execution_currency_code] } { output: [ sdk_order_execution_aggregates_mv.order_id, sdk_order_execution_aggregates_mv.filled_quantity, sdk_order_execution_aggregates_mv.total_amount, sdk_order_execution_aggregates_mv.total_fees, sdk_order_execution_aggregates_mv.average_price, sdk_order_execution_aggregates_mv.execution_currency_code ], stream key: [ sdk_order_execution_aggregates_mv.order_id ] }
├── Upstream { output: [ order_id, filled_quantity, total_amount, total_fees, average_price, execution_currency_code ], stream key: [] }
└── BatchPlanNode { output: [ order_id, filled_quantity, total_amount, total_fees, average_price, execution_currency_code ], stream key: [] }
Fragment 25092 (Actor 120528,120527)
StreamTableScan { table: assets_dm_next, columns: [id, name_en, name_ar, type, issue_currency_code, ticker, isin] } { output: [ assets_dm_next.id, assets_dm_next.name_en, assets_dm_next.name_ar, assets_dm_next.type, assets_dm_next.issue_currency_code, assets_dm_next.ticker, assets_dm_next.isin ], stream key: [ assets_dm_next.id ] }
├── Upstream { output: [ id, name_en, name_ar, type, issue_currency_code, ticker, isin ], stream key: [] }
└── BatchPlanNode { output: [ id, name_en, name_ar, type, issue_currency_code, ticker, isin ], stream key: [] }
Fragment 25093 (Actor 120564)
StreamProject { exprs: [SubtractWithTimeZone(now, '2 years':Interval, 'UTC':Varchar) as $expr1] } { output: [ $expr1 ], stream key: [] }
└── StreamNow { output: [ now ], stream key: [] }