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 优化策略、内存管理与溢出、容错恢复,以及性能调优实践
#目录
- Trino 在 coomia-dip 中的定位
- Coordinator-Worker 架构
- Connector 插件体系
- 查询计划与优化
- 跨引擎 Join 优化
- 动态过滤
- 内存管理与溢出
- 容错恢复
- 安全与隔离
- 性能调优
- 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 并行执行 |
| 下推 | 谓词、投影、聚合下推到数据源 |
| Join | CBO 自动选择策略,小表 Broadcast |
| 动态过滤 | 运行时将 Join 结果反馈为扫描过滤条件 |
| 内存 | 支持溢出到磁盘,防止 OOM |
| 容错 | Task 级重试,支持中间结果恢复 |
| 隔离 | 资源组 + 行级过滤实现 World 隔离 |
“下一篇预告:S8-20 将深入 nsjail 沙箱,探讨 coomia-dip 如何安全执行用户自定义函数。