返回博客

Trino 查询联邦深潜:跨引擎统一查询

1. [Trino 在 coomia-dip 中的定位](#1-trino-在-coomia-dip-中的定位)

Coomia发布于 2025年11月28日9 分钟阅读
分享本文Twitter / X

系列:S8 技术组件深潜 · 第 19 篇 | 难度:高级 | 阅读时间:20 分钟

Trino 查询联邦深潜:跨引擎统一查询

#TL;DR

  • Trino 是 coomia-dip Data Layer 的联邦查询引擎,实现跨 Iceberg(Lakehouse)、Doris(OLAP)、PostgreSQL(元数据)的统一 SQL 查询
  • 本文深入分析 Trino 的 Coordinator/Worker 架构、Connector 插件体系、查询计划优化(CBO/谓词下推/投影下推)、以及动态过滤与运行时自适应
  • 涵盖 coomia-dip 的 3 种 Catalog 配置、跨引擎 Join 优化策略、内存管理与溢出、容错恢复,以及性能调优实践

#目录

  1. Trino 在 coomia-dip 中的定位
  2. Coordinator-Worker 架构
  3. Connector 插件体系
  4. 查询计划与优化
  5. 跨引擎 Join 优化
  6. 动态过滤
  7. 内存管理与溢出
  8. 容错恢复
  9. 安全与隔离
  10. 性能调优
  11. Key Takeaways

#1. Trino 在 coomia-dip 中的定位

#1.1 联邦查询架构

Code
┌────────────────────────────────────────────────┐
│                 Trino Cluster                   │
│                                                  │
│  ┌──────────────────────────────────────────┐  │
│  │            Coordinator                    │  │
│  │  SQL Parser → Planner → Optimizer        │  │
│  └──────────────┬───────────────────────────┘  │
│                 │                               │
│      ┌──────────┼──────────┐                   │
│  ┌───┴───┐  ┌───┴───┐  ┌───┴───┐             │
│  │Worker1│  │Worker2│  │Worker3│             │
│  └───┬───┘  └───┬───┘  └───┬───┘             │
└──────┼──────────┼──────────┼──────────────────┘
       │          │          │
  ┌────┴───┐ ┌───┴────┐ ┌───┴──────┐
  │Iceberg │ │ Doris  │ │PostgreSQL│
  │Bronze/ │ │ OLAP   │ │ Metadata │
  │Silver  │ │ Gold   │ │          │
  └────────┘ └────────┘ └──────────┘

#1.2 coomia-dip Catalog 配置

PROPERTIES
# catalog/iceberg.properties
connector.name=iceberg
iceberg.catalog.type=nessie
iceberg.catalog.uri=http://nessie-server:19120/api/v2
iceberg.catalog.default-warehouse=s3://coomia-dip-lakehouse/warehouse

# catalog/doris.properties
connector.name=jdbc
connection-url=jdbc:mysql://doris-fe:9030/ontology
connection-user=coomia-dip
connection-password=${ENV:DORIS_PASSWORD}

# catalog/metadata.properties
connector.name=postgresql
connection-url=jdbc:postgresql://pg-metadata:5432/coomia-dip
connection-user=coomia-dip
connection-password=${ENV:PG_PASSWORD}

#2. Coordinator-Worker 架构

#2.1 查询执行流程

Code
Client
  │
  ▼ SQL
Coordinator
  │
  ├─ 1. SQL Parser (ANTLR)
  │     → AST (Abstract Syntax Tree)
  │
  ├─ 2. Analyzer
  │     → 解析表名/列名/类型
  │     → 连接 Catalog 获取元数据
  │
  ├─ 3. Planner
  │     → 生成逻辑计划 (LogicalPlan)
  │     → 规则优化 (RBO)
  │
  ├─ 4. Optimizer
  │     → 基于成本优化 (CBO)
  │     → 谓词下推 / 投影下推 / Join 重排
  │
  ├─ 5. Fragmenter
  │     → 将计划分割为 Fragment
  │     → 分配到 Worker
  │
  └─ 6. Scheduler
        → 调度 Fragment 执行
        → 收集结果

Worker
  │
  ├─ 执行 Fragment (Pipeline)
  ├─ 从 Connector 读取数据
  ├─ 内存中处理 (Filter/Project/Join/Aggregate)
  └─ 将结果发送给 Coordinator 或下游 Worker

#2.2 资源配置

PROPERTIES
# config.properties (Coordinator)
coordinator=true
node-scheduler.include-coordinator=false
http-server.http.port=8080
discovery.uri=http://coordinator:8080
query.max-memory=50GB
query.max-memory-per-node=10GB
query.max-total-memory-per-node=12GB

# config.properties (Worker)
coordinator=false
http-server.http.port=8080
discovery.uri=http://coordinator:8080
query.max-memory-per-node=10GB

#3. Connector 插件体系

#3.1 Connector SPI

Java
// Trino Connector 需要实现的核心接口
public interface ConnectorFactory {
    Connector create(String catalogName, Map<String, String> config);
}

public interface Connector {
    ConnectorMetadata getMetadata();      // Schema/Table 元数据
    ConnectorSplitManager getSplitManager();  // 数据分片
    ConnectorPageSourceProvider getPageSourceProvider();  // 数据读取
    ConnectorPageSinkProvider getPageSinkProvider();      // 数据写入
}

#3.2 coomia-dip 自定义 Connector 思路

Java
// 未来可以实现 OQL Connector,直接在 Trino 中执行 OQL
// connector.name=ontology
// ontology.grpc.host=control-Layer
// ontology.grpc.port=9090

// SELECT * FROM ontology.world1.Employee WHERE department = 'Engineering'
// → 转换为 gRPC 调用 → Control Layer → Data Layer

#4. 查询计划与优化

#4.1 谓词下推

SQL
-- 原始查询
SELECT e.name, d.department_name
FROM iceberg.ontology.employee e
JOIN doris.ontology.department d ON e.department_id = d.id
WHERE e.world_id = 'world-001'
  AND e.salary > 100000
  AND d.location = 'Shanghai'

-- 优化后:谓词下推到各数据源
-- Iceberg Scan: WHERE world_id = 'world-001' AND salary > 100000
-- Doris Scan: WHERE location = 'Shanghai'
-- Join 在 Trino Worker 内存中执行

#4.2 投影下推

SQL
-- 只选择需要的列,减少数据传输
-- 原始:SELECT * FROM large_table WHERE ...
-- 优化:SELECT col1, col2 FROM large_table WHERE ...
-- Connector 只读取 col1, col2 列(Iceberg 利用 Parquet 列式读取)

#4.3 基于成本的优化(CBO)

SQL
-- CBO 自动选择最优 Join 顺序
-- 小表 Join 大表 → 自动选择 Broadcast Join
-- 大表 Join 大表 → 自动选择 Distributed Hash Join

-- 查看执行计划
EXPLAIN (TYPE DISTRIBUTED)
SELECT e.name, d.department_name
FROM iceberg.ontology.employee e
JOIN doris.ontology.department d ON e.department_id = d.id;

-- 收集统计信息(用于 CBO)
ANALYZE iceberg.ontology.employee;

#5. 跨引擎 Join 优化

#5.1 Join 策略

策略适用场景原理
Broadcast Join一侧表小(< 100MB)小表广播到所有 Worker
Partitioned Join两侧都大按 Join Key 哈希分区
Cross Join笛卡尔积所有组合

#5.2 coomia-dip 的跨引擎 Join 示例

SQL
-- Iceberg (历史数据) JOIN Doris (实时聚合) JOIN PostgreSQL (元数据)
SELECT
    e.name,
    e.hire_date,
    d.department_name,
    m.total_sales
FROM iceberg.ontology.employee e           -- 历史全量
JOIN metadata.ontology.department d         -- 元数据
    ON e.department_id = d.id
JOIN doris.analytics.monthly_sales m       -- 实时聚合
    ON e.id = m.employee_id
WHERE e.world_id = 'world-001'
  AND m.month = '2026-03'
ORDER BY m.total_sales DESC
LIMIT 100;

#5.3 Join 优化提示

SQL
-- 手动指定 Join 分发策略
SELECT /*+ BROADCAST(d) */
    e.name, d.department_name
FROM iceberg.ontology.employee e
JOIN doris.ontology.department d
    ON e.department_id = d.id;

-- 手动指定 Join 类型
SELECT /*+ HASH_JOIN(e, d) */
    e.name, d.department_name
FROM iceberg.ontology.employee e
JOIN doris.ontology.department d
    ON e.department_id = d.id;

#6. 动态过滤

#6.1 运行时动态过滤

SQL
-- 动态过滤自动优化
SELECT e.name
FROM iceberg.ontology.employee e
JOIN doris.ontology.department d
    ON e.department_id = d.id
WHERE d.location = 'Shanghai';

-- Trino 执行流程:
-- 1. 先扫描 Doris department 表(小表),得到 Shanghai 的 department_id 列表
-- 2. 将 department_id 列表作为动态过滤条件
-- 3. 扫描 Iceberg employee 表时,跳过不匹配的分区和文件
-- → 大幅减少 Iceberg 扫描量

#6.2 配置

PROPERTIES
# 动态过滤配置
enable-dynamic-filtering=true
dynamic-filtering-max-per-driver-row-count=1000000
dynamic-filtering-max-per-driver-size=10MB
dynamic-filtering.large-broadcast.max-distinct-values-per-driver=10000

#7. 内存管理与溢出

#7.1 内存池

PROPERTIES
# 内存配置
query.max-memory=50GB                    # 集群级查询内存上限
query.max-memory-per-node=10GB          # 节点级查询内存上限
query.max-total-memory-per-node=12GB    # 节点总内存(含系统)
memory.heap-headroom-per-node=2GB       # JVM 堆预留

# 溢出配置
spill-enabled=true
spiller-spill-path=/data/trino/spill
spiller-max-used-space-threshold=0.9
spiller-threads=4

#7.2 溢出策略

Code
内存不足时的处理:

1. 尝试减少内存使用(压缩中间结果)
2. 触发溢出到磁盘(Spill)
3. 如果溢出后仍不够 → 查询失败

溢出支持的操作:
  ✅ Hash Join (Build 侧)
  ✅ Hash Aggregation
  ✅ Sort / Order By
  ✅ Window Function
  ❌ Cross Join
  ❌ Distinct

#8. 容错恢复

#8.1 Task Level 重试

PROPERTIES
# 容错配置 (Trino 400+)
retry-policy=TASK
task-retry-attempts-per-task=2
retry-initial-delay=10s
retry-max-delay=1m
retry-delay-scale-factor=2.0

# 交换溢出(用于容错重试时恢复中间数据)
exchange.deduplication-buffer-size=32MB
fault-tolerant-execution-target-task-input-size=256MB

#8.2 查询级重试

Python
# coomia-dip 查询重试包装
class TrinoQueryClient:
    async def execute_with_retry(self, sql: str, max_retries: int = 3) -> list[dict]:
        for attempt in range(max_retries):
            try:
                return await self._execute(sql)
            except TrinoExternalError as e:
                if attempt < max_retries - 1 and e.error_name in RETRIABLE_ERRORS:
                    await asyncio.sleep(2 ** attempt)
                    continue
                raise

RETRIABLE_ERRORS = {
    "GENERIC_INTERNAL_ERROR",
    "REMOTE_TASK_ERROR",
    "WORKER_NOT_FOUND",
    "EXCEEDED_TIME_LIMIT",
}

#9. 安全与隔离

#9.1 coomia-dip 的多 World 隔离

PROPERTIES
# 使用 Trino 的 System Access Control
access-control.name=file
security.config-file=/etc/trino/access-control.json
JSON
{
  "catalogs": [
    {
      "catalog": "iceberg",
      "allow": true
    }
  ],
  "schemas": [
    {
      "catalog": "iceberg",
      "schema": "world_.*",
      "owner": true
    }
  ],
  "tables": [
    {
      "catalog": "iceberg",
      "schema": "world_(.*)",
      "table": ".*",
      "privileges": ["SELECT"],
      "filter": "world_id = '${USER}'"
    }
  ]
}

#9.2 资源组

JSON
{
  "rootGroups": [
    {
      "name": "interactive",
      "maxQueued": 100,
      "hardConcurrencyLimit": 20,
      "softMemoryLimit": "40%",
      "schedulingPolicy": "fair"
    },
    {
      "name": "batch",
      "maxQueued": 500,
      "hardConcurrencyLimit": 10,
      "softMemoryLimit": "60%",
      "schedulingPolicy": "weighted_fair"
    }
  ]
}

#10. 性能调优

#10.1 关键配置

PROPERTIES
# Join 优化
join-distribution-type=AUTOMATIC
join-reordering-strategy=AUTOMATIC

# 并行度
task.concurrency=16
task.writer-count=4

# 节点调度
node-scheduler.min-candidates=10
node-scheduler.max-splits-per-node=100

# Iceberg 优化
iceberg.split-size=128MB
iceberg.max-partitions-per-writer=100

#10.2 查询分析

SQL
-- 查看详细执行统计
EXPLAIN ANALYZE
SELECT e.name, d.department_name
FROM iceberg.ontology.employee e
JOIN doris.ontology.department d ON e.department_id = d.id
WHERE e.world_id = 'world-001';

-- 关键指标
-- Input rows / Output rows — 过滤效果
-- Physical input / bytes — 实际读取量
-- Wall time — 执行耗时
-- Peak memory — 内存峰值

#11. Key Takeaways

主题关键结论
定位跨 Iceberg/Doris/PostgreSQL 的联邦查询
架构Coordinator 优化计划 + Worker 并行执行
下推谓词、投影、聚合下推到数据源
JoinCBO 自动选择策略,小表 Broadcast
动态过滤运行时将 Join 结果反馈为扫描过滤条件
内存支持溢出到磁盘,防止 OOM
容错Task 级重试,支持中间结果恢复
隔离资源组 + 行级过滤实现 World 隔离

下一篇预告:S8-20 将深入 nsjail 沙箱,探讨 coomia-dip 如何安全执行用户自定义函数。