Job is idle — throughput ~0; structure shown.
Fragment 25940 (Actor 130901,130902)
StreamMaterialize { columns: [party_id, account_group_type, top_allocations], stream_key: [party_id, account_group_type], pk_columns: [party_id, account_group_type], pk_conflict: NoCheck }
├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, jsonb_agg($expr1 order_by(row_number ASC)) ]
├── stream key: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type ]
└── StreamProject { exprs: [party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, jsonb_agg($expr1 order_by(row_number ASC))] }
├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, jsonb_agg($expr1 order_by(row_number ASC)) ]
├── stream key: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type ]
└── StreamHashAgg { group_key: [party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type], aggs: [jsonb_agg($expr1 order_by(row_number ASC)), count] }
├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, jsonb_agg($expr1 order_by(row_number ASC)), count ]
├── stream key: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type ]
└── MergeExecutor
├── output:
│ ┌── party_to_account_groups_mv_next.party_id
│ ├── party_to_account_groups_mv_next.type
│ ├── $expr1
│ ├── row_number
│ ├── party_to_account_groups_mv_next.parties.id
│ ├── party_to_account_groups_mv_next.null:Varchar
│ ├── party_to_account_groups_mv_next.$src
│ ├── intraday_position_by_distribution_mv_next.currency_code
│ ├── intraday_position_by_distribution_mv_next.position_type
│ ├── intraday_position_by_distribution_mv_next.distribution_type
│ ├── intraday_position_by_distribution_mv_next.taxonomy_node_id
│ └── party_to_account_groups_mv_next.account_group_id
└── stream key:
┌── party_to_account_groups_mv_next.parties.id
├── party_to_account_groups_mv_next.null:Varchar
├── party_to_account_groups_mv_next.$src
├── intraday_position_by_distribution_mv_next.currency_code
├── intraday_position_by_distribution_mv_next.position_type
├── intraday_position_by_distribution_mv_next.distribution_type
├── intraday_position_by_distribution_mv_next.taxonomy_node_id
└── party_to_account_groups_mv_next.account_group_id
Fragment 25941 (Actor 130903,130904)
StreamProject { exprs: [party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, JsonbBuildObject('assetClass':Varchar, intraday_position_by_distribution_mv_next.taxonomy_node_id, 'taxonomyNodeId':Varchar, intraday_position_by_distribution_mv_next.taxonomy_node_id, 'marketValue':Varchar, JsonbBuildObject('amount':Varchar, intraday_position_by_distribution_mv_next.market_value::Varchar, 'currencyCode':Varchar, intraday_position_by_distribution_mv_next.currency_code), 'fairValue':Varchar, JsonbBuildObject('amount':Varchar, intraday_position_by_distribution_mv_next.fair_value::Varchar, 'currencyCode':Varchar, intraday_position_by_distribution_mv_next.currency_code), 'weight':Varchar, intraday_position_by_distribution_mv_next.weight, 'rank':Varchar, row_number) as $expr1, row_number, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id, party_to_account_groups_mv_next.account_group_id] }
├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, $expr1, row_number, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id, party_to_account_groups_mv_next.account_group_id ]
├── stream key: [ party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id, party_to_account_groups_mv_next.account_group_id ]
└── MergeExecutor
├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, row_number, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src, party_to_account_groups_mv_next.account_group_id, intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ]
└── stream key: [ party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id, party_to_account_groups_mv_next.account_group_id ]
Fragment 25942 (Actor 130905,130906)
StreamSyncLogStore
├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, row_number, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src, party_to_account_groups_mv_next.account_group_id, intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ]
├── stream key: [ party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id, party_to_account_groups_mv_next.account_group_id ]
└── StreamHashJoin { type: Inner, predicate: party_to_account_groups_mv_next.account_group_id = intraday_position_by_distribution_mv_next.account_group_id }
├── output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.type, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, row_number, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src, party_to_account_groups_mv_next.account_group_id, intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ]
├── stream key: [ party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id, party_to_account_groups_mv_next.account_group_id ]
├── MergeExecutor { output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.account_group_id, party_to_account_groups_mv_next.type, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src ], stream key: [ party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src ] }
└── StreamOverWindow { window_functions: [row_number() OVER(PARTITION BY intraday_position_by_distribution_mv_next.account_group_id ORDER BY intraday_position_by_distribution_mv_next.market_value DESC, intraday_position_by_distribution_mv_next.taxonomy_node_id ASC ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)] } { output: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, row_number ], stream key: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id ] }
└── StreamGroupTopN { order: [intraday_position_by_distribution_mv_next.market_value DESC, intraday_position_by_distribution_mv_next.taxonomy_node_id ASC], limit: 5, offset: 0, group_key: [intraday_position_by_distribution_mv_next.account_group_id] } { output: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ], stream key: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id ] }
└── MergeExecutor { output: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ], stream key: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id ] }
Fragment 25943 (Actor 130908,130907)
StreamTableScan { table: party_to_account_groups_mv_next, columns: [party_id, account_group_id, type, parties.id, null:Varchar, $src] } { output: [ party_to_account_groups_mv_next.party_id, party_to_account_groups_mv_next.account_group_id, party_to_account_groups_mv_next.type, party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src ], stream key: [ party_to_account_groups_mv_next.parties.id, party_to_account_groups_mv_next.null:Varchar, party_to_account_groups_mv_next.$src ] }
├── Upstream { output: [ party_id, account_group_id, type, parties.id, null:Varchar, $src ], stream key: [] }
└── BatchPlanNode { output: [ party_id, account_group_id, type, parties.id, null:Varchar, $src ], stream key: [] }
Fragment 25944 (Actor 130900,130899)
StreamProject { exprs: [intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type] }
├── output: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ]
├── stream key: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id ]
└── StreamFilter { predicate: (intraday_position_by_distribution_mv_next.distribution_type = 'asset_classes':Varchar) AND (intraday_position_by_distribution_mv_next.position_type = 'POSITION':Varchar) } { output: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ], stream key: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id ] }
└── StreamTableScan { table: intraday_position_by_distribution_mv_next, columns: [account_group_id, taxonomy_node_id, currency_code, market_value, fair_value, weight, position_type, distribution_type] } { output: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.taxonomy_node_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.market_value, intraday_position_by_distribution_mv_next.fair_value, intraday_position_by_distribution_mv_next.weight, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type ], stream key: [ intraday_position_by_distribution_mv_next.account_group_id, intraday_position_by_distribution_mv_next.currency_code, intraday_position_by_distribution_mv_next.position_type, intraday_position_by_distribution_mv_next.distribution_type, intraday_position_by_distribution_mv_next.taxonomy_node_id ] }
├── Upstream { output: [ account_group_id, taxonomy_node_id, currency_code, market_value, fair_value, weight, position_type, distribution_type ], stream key: [] }
└── BatchPlanNode { output: [ account_group_id, taxonomy_node_id, currency_code, market_value, fair_value, weight, position_type, distribution_type ], stream key: [] }