StarRocks Embraces Paimon 2.0: Unifying Multimodal Analytics and Search in One SQL Engine
StarRocks now fully supports Paimon 2.0, integrating columnar analytics with vector, full-text, and scalar indexes into a single SQL execution framework for multimodal lakehouse workloads, enabling unified data management, retrieval, and relational analysis on object storage without data movement.
StarRocks has extended its lakehouse support to fully integrate with Paimon 2.0, creating a unified SQL execution architecture that combines traditional columnar OLAP scans with index-based search for multimodal data. Since proposing "From OLAP to Lakehouse" in 2023, StarRocks has continuously improved support for open table formats like Delta Lake, Iceberg, and Paimon, allowing interactive queries, multi-table joins, and business analytics directly on data lake storage.
Multimodal Lakehouse Challenges
Modern lakehouses must handle not only structured business records but also video, images, audio, documents, and sensor data. Models extract tags, text descriptions, and vectors from raw content, enabling new workflows: content-based sample retrieval, combined analysis with business records, and curated dataset preparation for model training and evaluation. This requires the lakehouse to manage diverse data types while simultaneously supporting scanning, retrieval, join analysis, and training data preparation.
Real-world examples illustrate the scale and complexity. The Waymo Open Motion Dataset (October 2025) contains 103,354 driving scenes of 20 seconds each, linking trajectories, maps, and image features. Google DeepMind's Open X-Embodiment (2023) aggregates over 1 million task episodes from 22 robot types. In these datasets, a single observation cannot be understood without temporal context, environment, and task goals; sample management must maintain cross-content associations and versioning.
Industry Use Cases
Autonomous Driving: From a rainy-night driving anomaly, retrieve similar scenes, filter by vehicle model and software version, correlate road conditions, point clouds, and vehicle state to identify systemic issues. Results form verifiable, re-labelable sample sets for model evaluation and training augmentation.
Embodied AI: From a robot insertion failure, retrieve similar visual states and action trajectories, combine with task type, policy version, contact forces, and execution results to diagnose root causes. Records become replayable, comparable task samples for policy evaluation and data recollection.
Advertising: From high-converting ads, retrieve creatives with similar composition, products, or style, then join with region, campaign cycle, and impression/click/conversion metrics to analyze content feature performance for creative iteration and A/B testing.
Gaming: Retrieve character and scene assets similar to a reference image or text description, then correlate with project, asset version, licensing, and release records to determine reusability, forming a traceable asset inventory.
These scenarios share common data problems: large-volume, long-retention raw media requiring cost-effective storage with consistency, and continuous association of raw content with evolving tags, features, vectors, and business metrics.
Paimon 2.0 Core Capabilities
Paimon 2.0 builds on its real-time lake table foundation to connect raw data management, feature engineering, and index retrieval, providing a unified data foundation for continuous multimodal data processing and consumption. The article focuses on multimodal Append tables with Row Tracking and Data Evolution enabled.
Driving Frame Example: From Raw Media to Searchable Features
A driving frame record initially stores capture time, vehicle model as regular columns, raw media in separate BLOB storage, and receives a row ID via Row Tracking. After model inference, the same row ID is used to backfill labels, evaluation scores, and vectors via Data Evolution. Labels and scores go to regular columns; vectors use dedicated VECTOR files. Unchanged business attributes and media files are reused. Through the shared row ID, initial writes and subsequent feature backfills always map to the same logical record.
Scalar and vector Global Indexes are built asynchronously on labels and vectors. Only indexes valid for the current query snapshot participate. After snapshot selection, Paimon uses visible data files and valid indexes, combining columns on demand via row ID. Raw content and backfilled features remain in one table, allowing StarRocks to filter, search, join, and aggregate via SQL.
StarRocks + Paimon 2.0: Unifying Analysis and Search in One SQL Framework
StarRocks' accumulated metadata management, SQL planning, distributed execution, and data caching across Delta Lake, Iceberg, and Paimon apply equally to multimodal lakehouses. Business data still needs joins, aggregations, and sorts; object storage access still requires I/O control; concurrent queries still need resource management. Paimon 2.0's Global Index adds a new index access path alongside traditional column scans.
Through coordination, StarRocks can combine scalar filtering, vector and full-text search, BLOB reads, and subsequent relational computation in a single SQL query. For example, a query can first narrow candidates by time, vehicle model, and labels, then execute similarity or keyword search; matched records can further join vehicle info, versions, and business metrics for aggregation and statistical analysis. All data stays in Paimon tables, retrieval and analysis use a consistent table snapshot, and a single query execution framework handles everything.
Global Indexes return row ID sets for scalar conditions, combined via AND/OR logic. Vector and full-text indexes also return candidate scores for merging and ranking. Scalar index candidate sets can serve as pre-filters for vector search. StarRocks chooses between column scans and index access based on query conditions and execution cost, unifying candidate merging, data fetch, and subsequent relational computation.
As illustrated, columnar scan and index access are two data access paths for different query patterns, but both ultimately enter StarRocks' relational computation framework for Filter, Join, aggregation, sort, and TopK operations. Paimon manages table versions, row identity, column files, and indexes; StarRocks handles access path planning, distributed execution, cache scheduling, and resource management. In one query, scans and index access can combine; predicate and TopK pushdown, data materialization timing, and operator ordering are all organized by the unified query plan.
From OLAP to Search: How Data Access Paths Change
OLAP: Metadata Pruning → Parallel Column Scan & Filter → Join, Aggregation, Sort
Traditional lake OLAP queries use columnar scan as entry. Queries leverage partition info and file-level min/max statistics to skip files that cannot satisfy conditions; Paimon readers further prune using row group and page-level statistics. This data skipping reduces scan volume based on range summaries but cannot pinpoint final matching records. Remaining data must still be read, decoded, and row-filtered. StarRocks organizes parallel read tasks by file or intra-file ranges, reads required columns in batches, decodes and filters in bulk, then feeds results into join, aggregation, and sort operators. This path suits reporting, scenario statistics, and large-range join analysis, optimizing for reduced actual scan volume, higher batch read/compute throughput, and balanced load across nodes.
Search: Index Search → Candidate Merge/Rank → Row ID Lookup → On-Demand Fetch & Subsequent Computation
Search uses indexes as the data access entry, directly finding candidate records via scalar conditions, keywords, or vector similarity. Scalar indexes typically return row ID sets; vector and full-text indexes also return distance or relevance scores. StarRocks merges, merges, and ranks candidates from different conditions, then reads required columns for the selected row IDs to continue filter, join, aggregation, and other relational computation.
During fetch, Paimon locates data files via row ID metadata in the query snapshot; readers fetch needed columns on demand. For Data Evolution tables, a single record's columns may span multiple files, requiring recombination based on row identity and data version. When candidate records are few, this path significantly reduces data read scope. However, index hit count does not fully represent fetch cost: if few row IDs scatter across many files, object storage requests may still be high. Actual cost depends on candidate count, hit file distribution, projected columns, and cache state.
StarRocks Execution Adaptations for Search
Extending from columnar scan to index search requires more than adding a read method. Index search, candidate merging, and data fetch have different parallelism boundaries and resource characteristics, requiring scheduler, materialization timing, and access path selection adaptations.
Coordinated Scheduling of Index Search and Data Fetch: StarRocks' existing distributed execution and data caching provide parallel scheduling foundations. Vector and full-text indexes split search tasks by index shard; returned row IDs locate files and columns for reading. Since index search and data fetch have different task boundaries, they need separate parallelism settings, scheduled with cache locality and node load awareness. Fetch phase can merge batch requests by target file, reusing metadata and data caches to reduce object storage requests and cross-node transfers.
Candidate Merging and Delayed Materialization: Using StarRocks' distributed merge and global delayed materialization, queries first exchange lightweight row IDs and scores, complete global candidate merge and ranking. Columns only needed for final output (media URLs, scene descriptions) can be batch-fetched after TopK results are determined, reducing useless data reads, decoding, and cross-node transfers. However, not all columns can be delayed: fields participating in filters or joins that affect candidate eligibility for final ranking must be computed before TopK; raw vectors for re-ranking may also need early reads. After fetch, data enters StarRocks' existing vectorized execution pipeline for joins and aggregations. For example, after selecting the top 100 most similar driving frames, further join vehicle info and aggregate by software version to analyze hit sample version distribution.
Cost-Based Access Path Selection: StarRocks' CBO uses statistics to estimate execution cost, combining predicate pushdown and join ordering. With Paimon index access, cost estimation must further cover index search scale, candidate count, hit file distribution, and Data Evolution column merge costs, unifying these with subsequent relational computation. Index access is not always faster: when candidates are too many or scattered across too many files, direct columnar scan may be more economical; for ANN vector queries, query latency and resource overhead must be compared under fixed recall targets. The goal is optimizer selection of the more appropriate execution path based on end-to-end cost.
Paimon 2.0 Best Practices: Data Management from Write to Query
Production experience distills multimodal table management into three areas: optimizing file count and layout, ordered backfill and organization of features, and timely index build and refresh. These respectively affect StarRocks' data pruning and column scans, on-demand reads and column recombination, and index search, candidate merge, and fetch costs. Query performance cannot be observed only at SQL execution; it requires holistic assessment of file layout, feature files, and index coverage to distinguish whether issues stem from query access paths or underlying data organization and maintenance.
Write and Compaction: File Layout Serving Queries
Design partitions by common query conditions, control file count by write volume: Time, business domain, and other common filters guide partition design to narrow read scope. Overly fine partitions or too-small write batches create many small files. For 1 TiB data, 1 MiB average file size yields ~1.05 million objects; 256 MiB yields ~4,096 objects. Drastically fewer files reduce metadata processing and object storage pressure. But files aren't better the larger; must balance write latency, task parallelism, pruning effectiveness, and read throughput.
Match compaction capacity to data growth rate: Accumulating small files increases manifest processing, file opens, and object storage requests. Continuously monitor file count, size distribution, and compaction backlog, correlating with StarRocks scan volume, object request count, and cache hit rate to judge if file organization truly improves query performance. Row Tracking tables use unaware Append mode, differing from primary key tables in file organization, so primary key bucket planning strategies don't directly apply. Configure compaction resources and frequency based on write speed, file growth, and query access patterns.
Data Evolution: Traceable, Maintainable Feature Backfill
Backfill by column, record feature provenance and version: After model upgrades, write only new or changed feature columns, reusing unchanged business attributes and raw media to avoid costly large-volume rewrites. Record model version, computation time, and processing status during backfill to ensure feature generation traceability. Vectors from different models or versions may not share feature space; query-time must define comparable scope. Row Tracking maintains row correspondence between derived features and original records; feature business meaning, generation method, and version relationships remain business-side responsibilities.
Include column file merge and index updates in backfill plans: Multiple partial updates reduce redundant writes of unchanged data but may increase query-time cost of combining multiple column files. Observe file counts involved in common query projections, and coordinate feature backfill, compaction, and index refresh to avoid read cost rising after write cost drops. When indexed feature columns change, verify index coverage consistency with current data version. For maintenance operations that may reassign row IDs, assess impact on existing indexes and external references.
Global Index: Configure by Condition, Scale, and Recall Targets
Different indexes solve different problems; selection cannot rely solely on field type. Configuration must consider query condition selectivity, data scale, index coverage, and retrieval precision/latency requirements.
Use index shards to balance parallelism and task overhead: Vector and full-text indexes control single build/search scale via index shards. Too many shards increase parallel task scheduling and candidate merge overhead; too large shards increase single search task compute and resource demand. BTree, Bitmap, and other indexes have their own data sorting and range organization; one shard strategy doesn't fit all. Index shards differ from table partitions or buckets: the former serves index build/search, the latter two determine table data organization and distribution; all three need separate planning.
Evaluate full query, verify index coverage: Index search is only part of the query chain. Performance evaluation must fix model version, distance metric, and recall target, separately observing cold/hot cache index search, candidate merge, and data fetch latency, combined with candidate count, hit file distribution, and projected columns to judge actual cost and configuration effectiveness.
DLF Managed Maintenance: Align Maintenance State with Query Performance
Alibaba Cloud Data Lake Formation (DLF) centrally manages Global Index async build, auto compaction, snapshot cleanup, partition lifecycle, and orphan file cleanup. Teams can correlate index coverage, maintenance task backlog, and resource consumption with StarRocks query latency, scan volume, and cache performance to tune maintenance resources and frequency.
Data write success and index query readiness are two distinct states to observe separately. DLF handles continuous lake table and index maintenance; StarRocks plans scans or index access per query version and completes subsequent relational computation. Both jointly affect time from write completion to efficient retrievability, so maintenance strategy must balance retrievability latency, query performance, and resource cost.
Future Work: StarRocks System Optimization Under Mixed Workloads
As Paimon multimodal tables carry more workloads, StarRocks must handle low-latency search and complex OLAP queries simultaneously, while contending with I/O competition from feature backfill, index build, and compaction. Future optimization focuses on workload isolation, index-aware query planning, and query execution coordination with lake table maintenance.
Isolate Mixed Workloads via StarRocks Compute Groups and Scheduling
Interactive search prioritizes stable tail latency; batch analytics prioritizes overall throughput. StarRocks can use separate compute groups for each query type, independently setting concurrency limits, resource budgets, and scaling policies. Further optimization will combine query characteristics to improve task scheduling, cache reuse, and hot data handling, reducing high-throughput task interference on online search.
Beyond compute resources, StarRocks data reads and background maintenance share object storage bandwidth. Front-end query read pressure must be evaluated alongside compaction, index build, and other background I/O consumption to reasonably allocate shared storage resources.
Make StarRocks Query Planning Index-Coverage Aware
A time gap typically exists between data write completion and index readiness. StarRocks must plan data access based on available indexes and their coverage in the query snapshot: where query semantics and reader capabilities allow, schedule supplemental reads for data not yet index-covered; if business chooses to search only indexed data, explicitly define the currently searchable data range. Future work will surface these differences in query plans and performance metrics, evaluating supplemental read overhead, letting businesses explicitly choose among result completeness, data freshness, and query latency.
Use StarRocks Query Metrics for Lake Table Maintenance Decisions
Continuous writes and feature backfill constantly generate new data files and column update files. Untimely compaction increases metadata processing and query-time column merge costs; overly frequent compaction consumes storage bandwidth and compute, disturbing existing caches. Therefore, StarRocks query latency, data read volume, and cache hit rate must guide scheduling of compaction, feature backfill, and index build, adjusting priorities by task backlog and online query pressure. For maintenance operations that may reassign row IDs or involve indexed column changes, index updates and query service switchover must be synchronized to ensure maintenance doesn't compromise query result completeness and consistency.
StarRocks has accumulated query optimization, distributed execution, data caching, and resource management capabilities in structured lakehouse queries. For multimodal data, these capabilities extend from traditional columnar analysis to index search, candidate merging, and on-demand fetch, letting different data access paths enter a single SQL compute framework. Paimon 2.0 further connects raw content, derived features, and indexes, enabling continuous multimodal data processing and evolution in the lake. StarRocks × Paimon collaboration is not merely adding search to lake tables, but enabling data management, content retrieval, and relational analysis on the same lake data. As analysis and retrieval converge, multimodal lakehouse becomes a long-term investment direction for the StarRocks community, with potential to become a key component in AI Data Infra connecting data management and query compute. The community welcomes more developers and users to bring real scenarios for discussion, validation, and contribution to advance this direction.
Signed-in readers can open the original source through BestHub's protected redirect.
This article has been distilled and summarized from source material, then republished for learning and reference. If you believe it infringes your rights, please contactand we will review it promptly.
StarRocks
StarRocks is an open‑source project under the Linux Foundation, focused on building a high‑performance, scalable analytical database that enables enterprises to create an efficient, unified lake‑house paradigm. It is widely used across many industries worldwide, helping numerous companies enhance their data analytics capabilities.
How this landed with the community
Was this worth your time?
0 Comments
Thoughtful readers leave field notes, pushback, and hard-won operational detail here.
