RWM Console cluster: risingwave-ai-hub-int.ai-hub-rwm.svc.cluster.local

← cluster insights objects benchmark_values_by_distribution_mv explain
Overview Objects Graph History
materialized view · insights.benchmark_values_by_distribution_mv profiled over 5s
seconds (1–30)

Job is idle — throughput ~0; structure shown.

Stateful hash join (4 state tables) — consider a temporal join for dimension lookupsAggregation state — unbounded unless keyed or temporally filtered
45 operators
Materialize · insights.benchmark_values_by_distribution_mv
0% idle 2 actors
Project
2 actors
HashAgg Aggregation state — unbounded unless keyed or temporally filtered
0% idle 2 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · benchmark_constituents_ft.benchmark_id = benchmarks_dm.id
2 actors
HashJoin · Inner · benchmark_constituents_ft.benchmark_id = benchmarks_dm.id Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · benchmarks_dm
2 actors
Filter · benchmarks_dm
0% idle 2 actors
StreamScan · benchmarks_dm
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
Project · (benchmark_constituents_ft.date >= asset_distributions_for_…
2 actors
Filter · (benchmark_constituents_ft.date >= asset_distributions_for_…
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
SyncLogStore · Inner · benchmark_constituents_ft.asset_id = asset_distributions_fo…
2 actors
HashJoin · Inner · benchmark_constituents_ft.asset_id = asset_distributions_fo… Stateful hash join (4 state tables) — consider a temporal join for dimension lookups
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · asset_distributions_for_consumers_mv
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Merge
2 actors
Exchange
0% idle 0 actors
LocalityProvider
0% idle 2 actors
Merge
2 actors
Exchange
0% idle 0 actors
StreamScan · benchmark_constituents_ft
0% idle 2 actors
BatchPlan
2 actors
Merge
2 actors
Heat = the operator's output-buffer backpressure over the sampling window. Click a node to fold its subtree.
Materialize · insights.benchmark_values_by_distribution_mv Materialize insights.benchmark_valu… idle · 2 actors Project Project — · 2 actors HashAgg HashAgg idle · 2 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project Project — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · benchmark_constituents_ft.benchmark_id = benchmarks_dm.id SyncLogStore Inner · benchmark_const… — · 2 actors HashJoin · Inner · benchmark_constituents_ft.benchmark_id = benchmarks_dm.id HashJoin Inner · benchmark_const… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · benchmarks_dm Project benchmarks_dm — · 2 actors Filter · benchmarks_dm Filter benchmarks_dm idle · 2 actors StreamScan · benchmarks_dm StreamScan benchmarks_dm idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors Project · (benchmark_constituents_ft.date >= asset_distributions_for_… Project (benchmark_constituents… — · 2 actors Filter · (benchmark_constituents_ft.date >= asset_distributions_for_… Filter (benchmark_constituents… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors SyncLogStore · Inner · benchmark_constituents_ft.asset_id = asset_distributions_fo… SyncLogStore Inner · benchmark_const… — · 2 actors HashJoin · Inner · benchmark_constituents_ft.asset_id = asset_distributions_fo… HashJoin Inner · benchmark_const… idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · asset_distributions_for_consumers_mv StreamScan asset_distributions_for… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors LocalityProvider LocalityProvider idle · 2 actors Merge Merge — · 2 actors Exchange Exchange idle · 0 actors StreamScan · benchmark_constituents_ft StreamScan benchmark_constituents_… idle · 2 actors BatchPlan BatchPlan — · 2 actors Merge Merge — · 2 actors
Streaming operator plan from EXPLAIN ANALYZE. Node heat = backpressure. Drag to pan, scroll to zoom.
Fragments (DESCRIBE FRAGMENTS) — click to expand
Fragment 14517 (Actor 159078,159077)
StreamMaterialize { columns: [benchmark_id, fact_date, distribution_type, taxonomy_node_id, taxonomy_code, weight], stream_key: [benchmark_id, fact_date, distribution_type, taxonomy_node_id, taxonomy_code], pk_columns: [benchmark_id, fact_date, distribution_type, taxonomy_node_id, taxonomy_code], pk_conflict: NoCheck }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, sum($expr1) ]
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code ]
└── StreamProject { exprs: [benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, sum($expr1)] }
    ├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, sum($expr1) ]
    ├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code ]
    └── StreamHashAgg { group_key: [benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code], aggs: [sum($expr1), count] }
        ├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, sum($expr1), count ]
        ├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code ]
        └── StreamLocalityProvider { locality_columns: [benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code] }
            ├── output:
            │   ┌── benchmark_constituents_ft.benchmark_id
            │   ├── benchmark_constituents_ft.date
            │   ├── asset_distributions_for_consumers_mv.distribution_type
            │   ├── asset_distributions_for_consumers_mv.taxonomy_node_id
            │   ├── asset_distributions_for_consumers_mv.taxonomy_code
            │   ├── $expr1
            │   ├── benchmark_constituents_ft.asset_id
            │   ├── asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id
            │   ├── asset_distributions_for_consumers_mv.asset_distributions_dm.dimension
            │   └── asset_distributions_for_consumers_mv.effective_start_date
            ├── stream key:
            │   ┌── benchmark_constituents_ft.benchmark_id
            │   ├── benchmark_constituents_ft.date
            │   ├── asset_distributions_for_consumers_mv.distribution_type
            │   ├── asset_distributions_for_consumers_mv.taxonomy_node_id
            │   ├── asset_distributions_for_consumers_mv.taxonomy_code
            │   ├── benchmark_constituents_ft.asset_id
            │   ├── asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id
            │   ├── asset_distributions_for_consumers_mv.asset_distributions_dm.dimension
            │   └── asset_distributions_for_consumers_mv.effective_start_date
            └── MergeExecutor
                ├── output:
                │   ┌── benchmark_constituents_ft.benchmark_id
                │   ├── benchmark_constituents_ft.date
                │   ├── asset_distributions_for_consumers_mv.distribution_type
                │   ├── asset_distributions_for_consumers_mv.taxonomy_node_id
                │   ├── asset_distributions_for_consumers_mv.taxonomy_code
                │   ├── $expr1
                │   ├── benchmark_constituents_ft.asset_id
                │   ├── asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id
                │   ├── asset_distributions_for_consumers_mv.asset_distributions_dm.dimension
                │   └── asset_distributions_for_consumers_mv.effective_start_date
                └── stream key:
                    ┌── benchmark_constituents_ft.benchmark_id
                    ├── benchmark_constituents_ft.asset_id
                    ├── benchmark_constituents_ft.date
                    ├── asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id
                    ├── asset_distributions_for_consumers_mv.taxonomy_node_id
                    ├── asset_distributions_for_consumers_mv.asset_distributions_dm.dimension
                    └── asset_distributions_for_consumers_mv.effective_start_date

Fragment 14518 (Actor 159085,159086)
StreamProject { exprs: [benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, (benchmark_constituents_ft.weight * asset_distributions_for_consumers_mv.share) as $expr1, benchmark_constituents_ft.asset_id, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date] }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, $expr1, benchmark_constituents_ft.asset_id, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
└── MergeExecutor
    ├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.weight, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, asset_distributions_for_consumers_mv.share, benchmark_constituents_ft.asset_id, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date, benchmarks_dm.id ]
    └── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]

Fragment 14519 (Actor 159084,159083)
StreamSyncLogStore
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.weight, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, asset_distributions_for_consumers_mv.share, benchmark_constituents_ft.asset_id, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date, benchmarks_dm.id ]
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
└── StreamHashJoin { type: Inner, predicate: benchmark_constituents_ft.benchmark_id = benchmarks_dm.id }
    ├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.weight, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, asset_distributions_for_consumers_mv.share, benchmark_constituents_ft.asset_id, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date, benchmarks_dm.id ]
    ├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
    ├── MergeExecutor
    │   ├── output:
    │   │   ┌── benchmark_constituents_ft.benchmark_id
    │   │   ├── benchmark_constituents_ft.date
    │   │   ├── benchmark_constituents_ft.weight
    │   │   ├── asset_distributions_for_consumers_mv.distribution_type
    │   │   ├── asset_distributions_for_consumers_mv.taxonomy_node_id
    │   │   ├── asset_distributions_for_consumers_mv.taxonomy_code
    │   │   ├── asset_distributions_for_consumers_mv.share
    │   │   ├── benchmark_constituents_ft.asset_id
    │   │   ├── asset_distributions_for_consumers_mv.asset_id
    │   │   ├── asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id
    │   │   ├── asset_distributions_for_consumers_mv.asset_distributions_dm.dimension
    │   │   └── asset_distributions_for_consumers_mv.effective_start_date
    │   └── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
    └── MergeExecutor { output: [ benchmarks_dm.id ], stream key: [ benchmarks_dm.id ] }

Fragment 14520 (Actor 159087,159088)
StreamLocalityProvider { locality_columns: [benchmark_constituents_ft.benchmark_id] }
├── output:
│   ┌── benchmark_constituents_ft.benchmark_id
│   ├── benchmark_constituents_ft.date
│   ├── benchmark_constituents_ft.weight
│   ├── asset_distributions_for_consumers_mv.distribution_type
│   ├── asset_distributions_for_consumers_mv.taxonomy_node_id
│   ├── asset_distributions_for_consumers_mv.taxonomy_code
│   ├── asset_distributions_for_consumers_mv.share
│   ├── benchmark_constituents_ft.asset_id
│   ├── asset_distributions_for_consumers_mv.asset_id
│   ├── asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id
│   ├── asset_distributions_for_consumers_mv.asset_distributions_dm.dimension
│   └── asset_distributions_for_consumers_mv.effective_start_date
├── stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
└── MergeExecutor
    ├── output:
    │   ┌── benchmark_constituents_ft.benchmark_id
    │   ├── benchmark_constituents_ft.date
    │   ├── benchmark_constituents_ft.weight
    │   ├── asset_distributions_for_consumers_mv.distribution_type
    │   ├── asset_distributions_for_consumers_mv.taxonomy_node_id
    │   ├── asset_distributions_for_consumers_mv.taxonomy_code
    │   ├── asset_distributions_for_consumers_mv.share
    │   ├── benchmark_constituents_ft.asset_id
    │   ├── asset_distributions_for_consumers_mv.asset_id
    │   ├── asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id
    │   ├── asset_distributions_for_consumers_mv.asset_distributions_dm.dimension
    │   └── asset_distributions_for_consumers_mv.effective_start_date
    └── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]

Fragment 14521 (Actor 159092,159091)
StreamProject { exprs: [benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.weight, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, asset_distributions_for_consumers_mv.share, benchmark_constituents_ft.asset_id, asset_distributions_for_consumers_mv.asset_id, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date] }
├── output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.weight, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, asset_distributions_for_consumers_mv.share, benchmark_constituents_ft.asset_id, asset_distributions_for_consumers_mv.asset_id, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
├── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
└── StreamFilter { predicate: (benchmark_constituents_ft.date >= asset_distributions_for_consumers_mv.effective_start_date) AND (benchmark_constituents_ft.date < asset_distributions_for_consumers_mv.effective_end_date) }
    ├── output:
    │   ┌── benchmark_constituents_ft.benchmark_id
    │   ├── benchmark_constituents_ft.date
    │   ├── benchmark_constituents_ft.asset_id
    │   ├── benchmark_constituents_ft.weight
    │   ├── asset_distributions_for_consumers_mv.asset_id
    │   ├── asset_distributions_for_consumers_mv.distribution_type
    │   ├── asset_distributions_for_consumers_mv.effective_start_date
    │   ├── asset_distributions_for_consumers_mv.effective_end_date
    │   ├── asset_distributions_for_consumers_mv.taxonomy_node_id
    │   ├── asset_distributions_for_consumers_mv.taxonomy_code
    │   ├── asset_distributions_for_consumers_mv.share
    │   ├── asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id
    │   └── asset_distributions_for_consumers_mv.asset_distributions_dm.dimension
    ├── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
    └── MergeExecutor
        ├── output:
        │   ┌── benchmark_constituents_ft.benchmark_id
        │   ├── benchmark_constituents_ft.date
        │   ├── benchmark_constituents_ft.asset_id
        │   ├── benchmark_constituents_ft.weight
        │   ├── asset_distributions_for_consumers_mv.asset_id
        │   ├── asset_distributions_for_consumers_mv.distribution_type
        │   ├── asset_distributions_for_consumers_mv.effective_start_date
        │   ├── asset_distributions_for_consumers_mv.effective_end_date
        │   ├── asset_distributions_for_consumers_mv.taxonomy_node_id
        │   ├── asset_distributions_for_consumers_mv.taxonomy_code
        │   ├── asset_distributions_for_consumers_mv.share
        │   ├── asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id
        │   └── asset_distributions_for_consumers_mv.asset_distributions_dm.dimension
        └── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]

Fragment 14522 (Actor 159090,159089)
StreamSyncLogStore
├── output:
│   ┌── benchmark_constituents_ft.benchmark_id
│   ├── benchmark_constituents_ft.date
│   ├── benchmark_constituents_ft.asset_id
│   ├── benchmark_constituents_ft.weight
│   ├── asset_distributions_for_consumers_mv.asset_id
│   ├── asset_distributions_for_consumers_mv.distribution_type
│   ├── asset_distributions_for_consumers_mv.effective_start_date
│   ├── asset_distributions_for_consumers_mv.effective_end_date
│   ├── asset_distributions_for_consumers_mv.taxonomy_node_id
│   ├── asset_distributions_for_consumers_mv.taxonomy_code
│   ├── asset_distributions_for_consumers_mv.share
│   ├── asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id
│   └── asset_distributions_for_consumers_mv.asset_distributions_dm.dimension
├── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
└── StreamHashJoin { type: Inner, predicate: benchmark_constituents_ft.asset_id = asset_distributions_for_consumers_mv.asset_id }
    ├── output:
    │   ┌── benchmark_constituents_ft.benchmark_id
    │   ├── benchmark_constituents_ft.date
    │   ├── benchmark_constituents_ft.asset_id
    │   ├── benchmark_constituents_ft.weight
    │   ├── asset_distributions_for_consumers_mv.asset_id
    │   ├── asset_distributions_for_consumers_mv.distribution_type
    │   ├── asset_distributions_for_consumers_mv.effective_start_date
    │   ├── asset_distributions_for_consumers_mv.effective_end_date
    │   ├── asset_distributions_for_consumers_mv.taxonomy_node_id
    │   ├── asset_distributions_for_consumers_mv.taxonomy_code
    │   ├── asset_distributions_for_consumers_mv.share
    │   ├── asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id
    │   └── asset_distributions_for_consumers_mv.asset_distributions_dm.dimension
    ├── stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
    ├── MergeExecutor { output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight ], stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date ] }
    └── MergeExecutor
        ├── output: [ asset_distributions_for_consumers_mv.asset_id, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.effective_start_date, asset_distributions_for_consumers_mv.effective_end_date, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, asset_distributions_for_consumers_mv.share, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension ]
        └── stream key: [ asset_distributions_for_consumers_mv.asset_id, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]

Fragment 14523 (Actor 159094,159093)
StreamLocalityProvider { locality_columns: [benchmark_constituents_ft.asset_id] } { output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight ], stream key: [ benchmark_constituents_ft.asset_id, benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date ] }
└── MergeExecutor { output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight ], stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id ] }

Fragment 14524 (Actor 159058,159057)
StreamTableScan { table: benchmark_constituents_ft, columns: [benchmark_id, date, asset_id, weight] } { output: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id, benchmark_constituents_ft.weight ], stream key: [ benchmark_constituents_ft.benchmark_id, benchmark_constituents_ft.date, benchmark_constituents_ft.asset_id ] }
├── Upstream { output: [ benchmark_id, date, asset_id, weight ], stream key: [] }
└── BatchPlanNode { output: [ benchmark_id, date, asset_id, weight ], stream key: [] }

Fragment 14525 (Actor 158036,158035)
StreamLocalityProvider { locality_columns: [asset_distributions_for_consumers_mv.asset_id] }
├── output: [ asset_distributions_for_consumers_mv.asset_id, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.effective_start_date, asset_distributions_for_consumers_mv.effective_end_date, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, asset_distributions_for_consumers_mv.share, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension ]
├── stream key: [ asset_distributions_for_consumers_mv.asset_id, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
└── MergeExecutor
    ├── output: [ asset_distributions_for_consumers_mv.asset_id, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.effective_start_date, asset_distributions_for_consumers_mv.effective_end_date, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, asset_distributions_for_consumers_mv.share, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension ]
    └── stream key: [ asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]

Fragment 14526 (Actor 158393,158392)
StreamTableScan { table: asset_distributions_for_consumers_mv, columns: [asset_id, distribution_type, effective_start_date, effective_end_date, taxonomy_node_id, taxonomy_code, share, taxonomy_nodes_dm.dimension_id, asset_distributions_dm.dimension] }
├── output: [ asset_distributions_for_consumers_mv.asset_id, asset_distributions_for_consumers_mv.distribution_type, asset_distributions_for_consumers_mv.effective_start_date, asset_distributions_for_consumers_mv.effective_end_date, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.taxonomy_code, asset_distributions_for_consumers_mv.share, asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension ]
├── stream key: [ asset_distributions_for_consumers_mv.taxonomy_nodes_dm.dimension_id, asset_distributions_for_consumers_mv.taxonomy_node_id, asset_distributions_for_consumers_mv.asset_id, asset_distributions_for_consumers_mv.asset_distributions_dm.dimension, asset_distributions_for_consumers_mv.effective_start_date ]
├── Upstream { output: [ asset_id, distribution_type, effective_start_date, effective_end_date, taxonomy_node_id, taxonomy_code, share, taxonomy_nodes_dm.dimension_id, asset_distributions_dm.dimension ], stream key: [] }
└── BatchPlanNode { output: [ asset_id, distribution_type, effective_start_date, effective_end_date, taxonomy_node_id, taxonomy_code, share, taxonomy_nodes_dm.dimension_id, asset_distributions_dm.dimension ], stream key: [] }

Fragment 14527 (Actor 158448,158449)
StreamProject { exprs: [benchmarks_dm.id] } { output: [ benchmarks_dm.id ], stream key: [ benchmarks_dm.id ] }
└── StreamFilter { predicate: IsNull(benchmarks_dm.disabled_at) } { output: [ benchmarks_dm.id, benchmarks_dm.disabled_at ], stream key: [ benchmarks_dm.id ] }
    └── StreamTableScan { table: benchmarks_dm, columns: [id, disabled_at] } { output: [ benchmarks_dm.id, benchmarks_dm.disabled_at ], stream key: [ benchmarks_dm.id ] }
        ├── Upstream { output: [ id, disabled_at ], stream key: [] }
        └── BatchPlanNode { output: [ id, disabled_at ], stream key: [] }