查询联邦:跨引擎统一查询
Tags: #QueryFederation #CrossEngine #Doris #DuckDB #Elasticsearch #智策平台
“系列:S3 数据基座 · 第 8 篇 | 难度:高级 | 阅读时间:20 分钟
查询联邦:跨引擎统一查询
Tags: #QueryFederation #CrossEngine #Doris #DuckDB #Elasticsearch #智策平台
#TL;DR
coomia-dip 平台的数据分布在 Doris(OLAP 分析)、DuckDB(嵌入式轻量查询)、Elasticsearch(全文搜索)、Iceberg(冷数据归档)等多个存储引擎中。用户不应关心数据存储位置——一条 OQL 查询应该自动路由到最合适的引擎并聚合结果。本文完整解析查询联邦(Query Federation)的架构设计、查询路由策略、跨引擎 JOIN 实现、数据传输优化和一致性保障机制。通过对标 Palantir Foundry 的统一数据访问层,展示如何在异构存储之上构建透明的联邦查询能力。
#1. 为什么需要查询联邦
#1.1 异构存储现实
coomia-dip 存储引擎分布:
┌─────────────┬────────────────────┬────────────────────┐
│ 引擎 │ 数据类型 │ 优势场景 │
├─────────────┼────────────────────┼────────────────────┤
│ Apache Doris │ 热数据、实时分析 │ 亚秒级 OLAP 查询 │
│ DuckDB │ 临时分析、小数据集 │ 零延迟嵌入式分析 │
│ Elasticsearch│ 全文索引、日志 │ 模糊搜索、分词匹配 │
│ Apache Iceberg│ 冷数据、历史归档 │ 时间旅行、Schema 演进 │
│ MinIO (S3) │ 对象存储、文件 │ 大文件、非结构化数据 │
└─────────────┴────────────────────┴────────────────────┘
问题:同一个 Entity 的属性可能分布在 3 个引擎中
- 基本属性 → Doris entity_common 表
- 全文描述 → Elasticsearch 倒排索引
- 历史版本 → Iceberg 快照
#1.2 Palantir Foundry 的方案
Foundry 的统一数据访问层:
┌──────────────────────────────────────────┐
│ Foundry API │
│ (用户只看到 Dataset 抽象) │
├──────────────────────────────────────────┤
│ Unified Query Engine │
│ (自动选择最优存储后端) │
├────────┬────────┬────────┬───────────────┤
│ Spark │ Trino │ ES │ Foundry SQL │
│ │ │ │ │
└────────┴────────┴────────┴───────────────┘
coomia-dip 的对应方案:
Foundry Dataset → Ontology Entity Type
Foundry API → OQL
Unified Engine → Query Federation Layer
Spark/Trino → Doris + DuckDB
#2. 联邦查询架构
#2.1 分层架构
查询联邦分层架构:
┌──────────────────────────────────────────┐
│ OQL Query Interface │
│ (用户提交 OQL 查询) │
├──────────────────────────────────────────┤
│ Query Planner │
│ (解析 → AST → 逻辑计划 → 物理计划) │
├──────────────────────────────────────────┤
│ Federation Coordinator │
│ (查询拆分 → 路由 → 执行 → 聚合) │
├──────────────────────────────────────────┤
│ Engine Adapters │
│ ┌────────┬────────┬────────┬──────────┐ │
│ │ Doris │ DuckDB │ ES │ Iceberg │ │
│ │Adapter │Adapter │Adapter │Adapter │ │
│ └────────┴────────┴────────┴──────────┘ │
├──────────────────────────────────────────┤
│ Data Transfer Layer │
│ (Arrow Flight 高性能数据传输) │
└──────────────────────────────────────────┘
#2.2 查询路由决策
class QueryRouter:
"""查询路由器:决定查询发送到哪个引擎"""
def __init__(self, catalog: DataCatalog):
self._catalog = catalog
def route(self, plan: LogicalPlan) -> RoutingDecision:
"""基于查询特征选择最优引擎"""
features = self._analyze_query_features(plan)
# 规则 1:全文搜索 → Elasticsearch
if features.has_full_text_search:
return RoutingDecision(
primary_engine=Engine.ELASTICSEARCH,
reason="Query contains full-text search predicates"
)
# 规则 2:时间旅行 → Iceberg (via Doris Catalog)
if features.has_time_travel:
return RoutingDecision(
primary_engine=Engine.ICEBERG,
reason="Query requires time-travel snapshots"
)
# 规则 3:小数据集 + 复杂计算 → DuckDB
if features.estimated_rows < 100_000 and features.complexity > 0.7:
return RoutingDecision(
primary_engine=Engine.DUCKDB,
reason="Small dataset with complex computation"
)
# 规则 4:大数据集 + OLAP → Doris
if features.estimated_rows >= 100_000:
return RoutingDecision(
primary_engine=Engine.DORIS,
reason="Large dataset OLAP query"
)
# 默认 → Doris
return RoutingDecision(
primary_engine=Engine.DORIS,
reason="Default routing"
)
def route_federated(self, plan: LogicalPlan) -> FederatedPlan:
"""处理需要跨引擎的查询"""
sub_plans = self._split_plan(plan)
assignments = {}
for sub_plan in sub_plans:
engine = self.route(sub_plan).primary_engine
assignments[sub_plan.id] = engine
return FederatedPlan(
sub_plans=sub_plans,
assignments=assignments,
join_strategy=self._choose_join_strategy(sub_plans)
)
#3. 引擎适配器
#3.1 适配器接口
from abc import ABC, abstractmethod
from typing import AsyncIterator
import pyarrow as pa
class EngineAdapter(ABC):
"""引擎适配器基类"""
@abstractmethod
async def execute(
self, plan: PhysicalPlan
) -> AsyncIterator[pa.RecordBatch]:
"""执行物理计划,返回 Arrow RecordBatch 流"""
...
@abstractmethod
def translate(self, plan: LogicalPlan) -> str:
"""将逻辑计划翻译为引擎原生查询语言"""
...
@abstractmethod
def estimate_cost(self, plan: LogicalPlan) -> QueryCost:
"""估算查询成本"""
...
@abstractmethod
def get_statistics(self, entity_type: str) -> TableStatistics:
"""获取表统计信息"""
...
#3.2 Doris 适配器
class DorisAdapter(EngineAdapter):
"""Apache Doris OLAP 引擎适配器"""
def __init__(self, connection_pool: DorisConnectionPool):
self._pool = connection_pool
def translate(self, plan: LogicalPlan) -> str:
compiler = DorisSQLCompiler()
return plan.accept(compiler)
async def execute(
self, plan: PhysicalPlan
) -> AsyncIterator[pa.RecordBatch]:
sql = self.translate(plan.logical_plan)
async with self._pool.acquire() as conn:
# 使用 Arrow Flight SQL 协议获取结果
flight_client = await conn.get_flight_client()
ticket = await flight_client.execute(sql)
async for batch in flight_client.do_get(ticket):
yield batch
def estimate_cost(self, plan: LogicalPlan) -> QueryCost:
# 使用 Doris EXPLAIN 获取成本估算
sql = self.translate(plan)
explain = self._pool.execute_sync(f"EXPLAIN {sql}")
return self._parse_explain(explain)
#3.3 DuckDB 适配器
class DuckDBAdapter(EngineAdapter):
"""DuckDB 嵌入式分析引擎适配器"""
def __init__(self, db_path: str = ":memory:"):
self._conn = duckdb.connect(db_path)
self._setup_iceberg_extension()
def _setup_iceberg_extension(self):
"""加载 Iceberg 扩展以读取 MinIO 上的 Iceberg 表"""
self._conn.execute("INSTALL iceberg; LOAD iceberg;")
self._conn.execute("INSTALL httpfs; LOAD httpfs;")
self._conn.execute(f"""
SET s3_endpoint='minio.internal:9000';
SET s3_access_key_id='minioadmin';
SET s3_secret_access_key='minioadmin';
SET s3_use_ssl=false;
SET s3_url_style='path';
""")
def translate(self, plan: LogicalPlan) -> str:
compiler = DuckDBSQLCompiler()
return plan.accept(compiler)
async def execute(
self, plan: PhysicalPlan
) -> AsyncIterator[pa.RecordBatch]:
sql = self.translate(plan.logical_plan)
result = self._conn.execute(sql)
while True:
batch = result.fetch_arrow_table()
if batch.num_rows == 0:
break
for record_batch in batch.to_batches(max_chunksize=8192):
yield record_batch
#3.4 Elasticsearch 适配器
class ElasticsearchAdapter(EngineAdapter):
"""Elasticsearch 全文搜索适配器"""
def __init__(self, es_client: AsyncElasticsearch):
self._es = es_client
def translate(self, plan: LogicalPlan) -> dict:
"""翻译为 Elasticsearch DSL"""
compiler = ESDSLCompiler()
return plan.accept(compiler)
async def execute(
self, plan: PhysicalPlan
) -> AsyncIterator[pa.RecordBatch]:
es_query = self.translate(plan.logical_plan)
index = self._resolve_index(plan.entity_type)
# 使用 scroll API 处理大结果集
async for hits in self._scroll_search(index, es_query):
batch = self._hits_to_arrow(hits)
yield batch
def _hits_to_arrow(self, hits: list[dict]) -> pa.RecordBatch:
"""将 ES hits 转换为 Arrow RecordBatch"""
columns = {}
for hit in hits:
source = hit['_source']
for key, value in source.items():
columns.setdefault(key, []).append(value)
arrays = [pa.array(values) for values in columns.values()]
names = list(columns.keys())
return pa.RecordBatch.from_arrays(arrays, names=names)
#4. 跨引擎 JOIN
#4.1 JOIN 策略
跨引擎 JOIN 策略选择:
┌─────────────────┬────────────────────┬───────────────────┐
│ 策略 │ 适用场景 │ 实现方式 │
├─────────────────┼────────────────────┼───────────────────┤
│ Broadcast Join │ 一侧数据量小 │ 小表广播到大表引擎 │
│ │ (<10MB) │ │
├─────────────────┼────────────────────┼───────────────────┤
│ Semi-Join Push │ 过滤后数据量小 │ 先执行过滤,将 ID │
│ │ │ 列表推送到对端 │
├─────────────────┼────────────────────┼───────────────────┤
│ Hash Join │ 双方数据量均中等 │ 拉取到协调器做 │
│ (at Coordinator) │ (<1M rows each) │ 本地 Hash Join │
├─────────────────┼────────────────────┼───────────────────┤
│ Sort-Merge Join │ 双方数据量大且有序 │ 流式归并,内存占用低 │
├─────────────────┼────────────────────┼───────────────────┤
│ Materialized │ 频繁跨引擎 JOIN │ 预计算物化到单引擎 │
│ View │ │ │
└─────────────────┴────────────────────┴───────────────────┘
#4.2 Semi-Join 推送优化
class SemiJoinPushdown:
"""Semi-Join 推送优化器"""
async def execute(
self,
left_adapter: EngineAdapter,
left_plan: PhysicalPlan,
right_adapter: EngineAdapter,
right_plan: PhysicalPlan,
join_key: str
) -> AsyncIterator[pa.RecordBatch]:
# Step 1: 在较小一侧执行查询,提取 JOIN key
small_side = await self._determine_small_side(
left_adapter, left_plan,
right_adapter, right_plan
)
if small_side == 'left':
keys = await self._extract_keys(left_adapter, left_plan, join_key)
# Step 2: 将 key 列表推送到右侧作为 IN 过滤
right_plan_filtered = self._inject_in_filter(
right_plan, join_key, keys
)
# Step 3: 在右侧执行过滤后的查询
right_results = right_adapter.execute(right_plan_filtered)
# Step 4: 在协调器做最终 JOIN
left_results = left_adapter.execute(left_plan)
async for batch in self._local_hash_join(
left_results, right_results, join_key
):
yield batch
#5. 联邦协调器
#5.1 查询执行流程
联邦查询执行流程:
1. 用户提交 OQL:
FETCH Person
SELECT name, age, description, historical_role
WHERE name LIKE '%张%' AND age > 30
AT TIME '2024-01-01'
2. 查询分析器识别数据源:
- name, age → Doris (entity_common)
- description → Elasticsearch (全文搜索)
- LIKE '%张%' → Elasticsearch (中文分词)
- AT TIME → Iceberg (时间旅行)
3. 查询拆分:
Sub-Query A (ES): 搜索 name LIKE '%张%' → 返回 entity_id 列表
Sub-Query B (Doris): SELECT name, age FROM entity_common
WHERE entity_id IN (...) AND age > 30
Sub-Query C (Iceberg): 读取 2024-01-01 时间快照的 historical_role
4. 执行计划:
[ES: 全文搜索] ──→ entity_id 列表
↓
[Doris: 属性查询] ←── IN 过滤 (Semi-Join Push)
↓
[Iceberg: 历史数据] ←── entity_id 列表
↓
[协调器: Hash Join] ──→ 最终结果
#5.2 并行执行引擎
class FederationCoordinator:
"""联邦查询协调器"""
def __init__(self, adapters: dict[Engine, EngineAdapter]):
self._adapters = adapters
async def execute(self, federated_plan: FederatedPlan) -> pa.Table:
# 构建执行 DAG
dag = self._build_execution_dag(federated_plan)
# 拓扑排序确定执行顺序
execution_order = dag.topological_sort()
results: dict[str, pa.Table] = {}
for level in execution_order:
# 同一层级的子查询可以并行执行
tasks = []
for sub_plan_id in level:
sub_plan = federated_plan.sub_plans[sub_plan_id]
engine = federated_plan.assignments[sub_plan_id]
adapter = self._adapters[engine]
# 注入上游结果作为参数
enriched_plan = self._inject_upstream_results(
sub_plan, results
)
tasks.append(self._execute_sub_plan(
adapter, enriched_plan, sub_plan_id
))
# 并行执行
level_results = await asyncio.gather(*tasks)
for sub_plan_id, result in zip(level, level_results):
results[sub_plan_id] = result
# 最终聚合
return self._merge_results(results, federated_plan)
async def _execute_sub_plan(
self,
adapter: EngineAdapter,
plan: PhysicalPlan,
plan_id: str
) -> pa.Table:
batches = []
async for batch in adapter.execute(plan):
batches.append(batch)
return pa.Table.from_batches(batches)
#6. 数据传输优化
#6.1 Arrow Flight 传输
跨引擎数据传输协议选择:
┌────────────────┬──────────┬──────────┬──────────┐
│ 协议 │ 吞吐量 │ 延迟 │ 序列化 │
├────────────────┼──────────┼──────────┼──────────┤
│ JDBC/ODBC │ ~100 MB/s│ 中 │ 行序列化 │
│ Arrow Flight │ ~2 GB/s │ 低 │ 零拷贝 │
│ Arrow Flight SQL│ ~1.5 GB/s│ 低 │ 零拷贝 │
│ gRPC Protobuf │ ~500 MB/s│ 低 │ Protobuf │
└────────────────┴──────────┴──────────┴──────────┘
选择:Arrow Flight SQL
- 与 Doris 原生集成
- 零拷贝内存映射
- 列式传输天然适合分析查询
#6.2 数据压缩与分批
class DataTransferOptimizer:
"""数据传输优化器"""
def __init__(self, config: TransferConfig):
self._config = config
async def transfer_with_compression(
self,
source: AsyncIterator[pa.RecordBatch],
target_engine: Engine
) -> AsyncIterator[pa.RecordBatch]:
"""带压缩的数据传输"""
buffer = []
buffer_size = 0
async for batch in source:
buffer.append(batch)
buffer_size += batch.nbytes
if buffer_size >= self._config.batch_size_bytes:
merged = pa.Table.from_batches(buffer)
# 根据目标引擎选择最优压缩
if target_engine == Engine.DORIS:
compressed = self._compress_lz4(merged)
elif target_engine == Engine.DUCKDB:
compressed = merged # DuckDB 本地,无需压缩
else:
compressed = self._compress_zstd(merged)
for out_batch in compressed.to_batches(
max_chunksize=self._config.chunk_size
):
yield out_batch
buffer = []
buffer_size = 0
# 处理剩余数据
if buffer:
merged = pa.Table.from_batches(buffer)
for out_batch in merged.to_batches(
max_chunksize=self._config.chunk_size
):
yield out_batch
#7. 一致性与事务保障
#7.1 快照一致性
联邦查询一致性模型:
问题:跨引擎查询期间,数据可能被并发修改
- t0: 从 Doris 读取 Person 列表
- t1: Person 数据被更新
- t2: 从 ES 读取 description
→ 结果不一致:Person 列表是 t0 版本,description 是 t2 版本
解决方案:Snapshot Isolation
┌─────────────────────────────────────────┐
│ Federation Snapshot Manager │
│ │
│ 1. 查询开始时获取全局快照时间戳 T_snap │
│ 2. 所有子查询使用 T_snap 作为读取时间点 │
│ │
│ Doris: 读取 T_snap 时刻的 MVCC 版本 │
│ Iceberg: 使用 T_snap 对应的 snapshot_id │
│ ES: 使用 Point-in-Time API │
│ DuckDB: 直接读取(嵌入式,无并发问题) │
└─────────────────────────────────────────┘
#7.2 快照管理
class SnapshotManager:
"""联邦查询快照管理器"""
async def create_snapshot(self) -> FederatedSnapshot:
"""创建跨引擎一致快照"""
timestamp = datetime.utcnow()
# 获取各引擎快照标识
doris_snapshot = await self._doris.get_snapshot_at(timestamp)
iceberg_snapshot = await self._iceberg.get_snapshot_at(timestamp)
es_pit = await self._es.open_point_in_time(
index="entity_*",
keep_alive="5m"
)
return FederatedSnapshot(
timestamp=timestamp,
doris_snapshot=doris_snapshot,
iceberg_snapshot_id=iceberg_snapshot.snapshot_id,
es_point_in_time=es_pit,
ttl=timedelta(minutes=5)
)
async def release_snapshot(self, snapshot: FederatedSnapshot):
"""释放快照资源"""
await self._es.close_point_in_time(snapshot.es_point_in_time)
#8. 查询缓存与物化
#8.1 结果缓存
class FederatedQueryCache:
"""联邦查询结果缓存"""
def __init__(self, redis_client: Redis, max_cache_size_mb: int = 512):
self._redis = redis_client
self._max_size = max_cache_size_mb * 1024 * 1024
async def get_or_execute(
self,
query_hash: str,
executor: Callable[[], Awaitable[pa.Table]]
) -> pa.Table:
# 检查缓存
cached = await self._redis.get(f"fed_cache:{query_hash}")
if cached:
return self._deserialize(cached)
# 执行查询
result = await executor()
# 缓存结果(如果不太大)
serialized = self._serialize(result)
if len(serialized) < self._max_size // 100: # 单条 < 总缓存的 1%
await self._redis.setex(
f"fed_cache:{query_hash}",
timedelta(minutes=5),
serialized
)
return result
#8.2 自动物化建议
物化视图自动推荐:
监控系统发现:
- 查询 Q1(Doris + ES JOIN)每天执行 500 次
- 平均耗时 2.3 秒
- 结果集 < 50MB
推荐:
CREATE MATERIALIZED VIEW mv_person_with_desc AS
SELECT p.entity_id, p.name, p.age, es.description
FROM entity_common p
JOIN es_entity_index es ON p.entity_id = es.entity_id
WHERE p.entity_type = 'Person'
REFRESH EVERY 5 MINUTES;
效果:
- 查询时间:2.3s → 50ms
- 跨引擎 JOIN → 单引擎查询
- 代价:额外 50MB 存储 + 每 5 分钟刷新
#9. 监控与可观测性
#9.1 联邦查询 Trace
分布式追踪示例:
Trace ID: fed-20240301-abc123
Total Duration: 1.2s
├── [0-50ms] Query Parse & Plan
│ ├── Lexer: 2ms
│ ├── Parser: 5ms
│ ├── Semantic Analysis: 15ms
│ └── Federation Planning: 28ms
│
├── [50-800ms] Sub-Query Execution (parallel)
│ ├── [50-200ms] ES: full-text search "张%" → 150 entity_ids
│ ├── [50-600ms] Doris: SELECT ... WHERE entity_id IN (...) → 120 rows
│ └── [50-800ms] Iceberg: time-travel read → 120 rows
│
├── [800-1100ms] Cross-Engine Join
│ ├── Hash Build: 50ms (ES + Iceberg results)
│ ├── Hash Probe: 200ms (Doris results)
│ └── Projection: 50ms
│
└── [1100-1200ms] Result Serialization & Transfer
└── Arrow Flight: 100ms (120 rows, 48KB)
#9.2 性能指标
联邦查询关键指标:
┌────────────────────┬──────────┬──────────┬──────────┐
│ 指标 │ P50 │ P95 │ P99 │
├────────────────────┼──────────┼──────────┼──────────┤
│ 单引擎查询延迟 │ 50ms │ 200ms │ 500ms │
│ 双引擎联邦延迟 │ 200ms │ 800ms │ 2s │
│ 三引擎联邦延迟 │ 500ms │ 1.5s │ 3s │
│ 跨引擎数据传输 │ 10ms │ 100ms │ 500ms │
│ JOIN 开销 │ 50ms │ 200ms │ 1s │
│ 缓存命中率 │ 65% │ - │ - │
└────────────────────┴──────────┴──────────┴──────────┘
#10. 测试策略
#10.1 引擎模拟测试
class TestQueryFederation:
@pytest.fixture
def mock_adapters(self):
"""模拟各引擎适配器"""
doris = MockDorisAdapter(data={
'Person': [
{'entity_id': '1', 'name': '张三', 'age': 35},
{'entity_id': '2', 'name': '李四', 'age': 28},
]
})
es = MockESAdapter(data={
'Person': [
{'entity_id': '1', 'description': '资深工程师'},
{'entity_id': '2', 'description': '产品经理'},
]
})
return {'doris': doris, 'es': es}
async def test_cross_engine_join(self, mock_adapters):
coordinator = FederationCoordinator(mock_adapters)
result = await coordinator.execute(
parse("FETCH Person SELECT name, description WHERE age > 30")
)
assert len(result) == 1
assert result[0]['name'] == '张三'
assert result[0]['description'] == '资深工程师'
async def test_routing_full_text_to_es(self, mock_adapters):
router = QueryRouter(catalog)
plan = parse("FETCH Person WHERE description LIKE '%工程师%'")
decision = router.route(plan)
assert decision.primary_engine == Engine.ELASTICSEARCH
#Key Takeaways
-
查询联邦让用户忽略存储异构性:一条 OQL 查询自动路由到 Doris、DuckDB、Elasticsearch、Iceberg 等引擎,用户无需知道数据物理存储位置。
-
查询路由策略是性能关键:基于数据量、查询复杂度、存储特性的自动路由,比手动指定引擎更优。
-
跨引擎 JOIN 需要精细的策略选择:Broadcast Join、Semi-Join Push、Hash Join、Sort-Merge Join 各有适用场景,Federation Coordinator 根据统计信息自动选择。
-
Arrow Flight 是跨引擎数据传输的最优协议:零拷贝列式传输,吞吐量是传统 JDBC 的 15-20 倍。
-
快照一致性是联邦查询的基础保障:跨引擎查询必须使用统一快照时间戳,否则结果不一致。
#Next Article
下一篇 S3-09《查询优化:从逻辑计划到物理执行》 将深入单引擎查询优化,包括谓词下推、列裁剪、分区剪枝和 Doris 物化视图自动匹配。
Tags: #QueryFederation #CrossEngine #Doris #DuckDB #Elasticsearch #Iceberg #ArrowFlight #SnapshotIsolation #FederatedJoin #智策平台 #coomia-dip #数据基座