Transform 执行器:多引擎适配层
Tags: #TransformExecutor #MultiEngine #Flink #Spark #DuckDB #智策平台
“系列:S3 数据基座 · 第 20 篇 | 难度:高级 | 阅读时间:20 分钟
Transform 执行器:多引擎适配层
Tags: #TransformExecutor #MultiEngine #Flink #Spark #DuckDB #智策平台
#TL;DR
Pipeline DSL 中的 Transform 步骤最终需要在具体计算引擎上执行。coomia-dip 的 Transform 执行器是一个多引擎适配层——相同的 Transform 定义可以编译为 Flink DataStream 算子、Spark DataFrame 操作或 DuckDB SQL 语句。本文完整解析 Transform 执行器的架构设计,包括统一 Transform 接口定义、Flink/Spark/DuckDB 三引擎的适配实现、引擎自动选择策略、Schema 跨引擎一致性保障、性能基准对比和故障切换机制。
#1. 为什么需要多引擎适配
#1.1 不同引擎的优势场景
三大引擎优势对比:
┌─────────────────┬──────────────┬──────────────┬──────────────┐
│ 场景 │ Flink │ Spark │ DuckDB │
├─────────────────┼──────────────┼──────────────┼──────────────┤
│ 实时流处理 │ ★★★★★ │ ★★★ │ ★ │
│ CDC 数据摄入 │ ★★★★★ │ ★★ │ ★ │
│ 大规模批处理 │ ★★★ │ ★★★★★ │ ★★ │
│ 交互式分析 │ ★★ │ ★★★ │ ★★★★★ │
│ 小数据集处理 │ ★★ │ ★★ │ ★★★★★ │
│ 窗口聚合 │ ★★★★★ │ ★★★★ │ ★★ │
│ 嵌入式执行 │ ✗ │ ✗ │ ★★★★★ │
│ 资源占用 │ 高 │ 高 │ 极低 │
└─────────────────┴──────────────┴──────────────┴──────────────┘
#1.2 统一接口的价值
统一 Transform 接口的价值:
用户视角:
pipeline.transform(filter(lambda r: r["age"] > 30))
→ 不关心底层是 Flink/Spark/DuckDB
系统视角:
→ 根据数据量、延迟要求和资源状况自动选择引擎
→ 一个引擎故障时自动切换到另一个
→ 开发/测试用 DuckDB(秒级启动),生产用 Flink
#2. Transform 接口定义
#2.1 统一接口
from abc import ABC, abstractmethod
from typing import Generic, TypeVar
import pyarrow as pa
T = TypeVar('T')
class TransformOperator(ABC, Generic[T]):
"""统一 Transform 算子接口"""
@property
@abstractmethod
def name(self) -> str: ...
@abstractmethod
def input_schema(self) -> pa.Schema: ...
@abstractmethod
def output_schema(self) -> pa.Schema: ...
@abstractmethod
def compile_flink(self) -> 'FlinkOperator': ...
@abstractmethod
def compile_spark(self) -> 'SparkTransform': ...
@abstractmethod
def compile_duckdb(self) -> str: ...
class FilterOperator(TransformOperator):
"""过滤算子"""
def __init__(self, predicate: str | Callable):
self._predicate = predicate
def output_schema(self) -> pa.Schema:
return self.input_schema() # Filter 不改变 Schema
def compile_flink(self) -> 'FlinkFilterFunction':
if callable(self._predicate):
return FlinkLambdaFilter(self._predicate)
return FlinkExprFilter(self._predicate)
def compile_spark(self) -> Column:
if callable(self._predicate):
return SparkUDFFilter(self._predicate)
return F.expr(self._predicate)
def compile_duckdb(self) -> str:
if callable(self._predicate):
raise CompilationError(
"Lambda predicates not supported in DuckDB mode. "
"Use SQL expression string instead."
)
return f"WHERE {self._predicate}"
class MapOperator(TransformOperator):
"""映射算子"""
def __init__(self, function: Callable[[dict], dict]):
self._function = function
self._output_fields: list[FieldDef] | None = None
def compile_flink(self) -> 'FlinkMapFunction':
return FlinkLambdaMap(self._function)
def compile_spark(self) -> 'SparkMapTransform':
return SparkUDFMap(self._function)
def compile_duckdb(self) -> str:
# 静态分析 lambda 函数,提取字段映射
field_mapping = self._analyze_lambda(self._function)
select_parts = [
f"{expr} AS {alias}"
for alias, expr in field_mapping.items()
]
return f"SELECT {', '.join(select_parts)}"
class AggregateOperator(TransformOperator):
"""聚合算子"""
def __init__(self, group_by: list[str], aggregations: dict):
self._group_by = group_by
self._aggregations = aggregations
def compile_flink(self) -> 'FlinkAggregation':
return FlinkKeyedAggregation(
key_selector=self._group_by,
aggregations=self._aggregations
)
def compile_spark(self) -> 'SparkGroupedAgg':
return SparkGroupByAgg(
group_cols=self._group_by,
agg_exprs=self._aggregations
)
def compile_duckdb(self) -> str:
agg_parts = []
for alias, (col, func) in self._aggregations.items():
agg_parts.append(f"{func.upper()}({col}) AS {alias}")
group_cols = ', '.join(self._group_by)
agg_cols = ', '.join(agg_parts)
return (
f"SELECT {group_cols}, {agg_cols} "
f"GROUP BY {group_cols}"
)
#3. Flink 适配器
#3.1 DataStream 编译
class FlinkTransformCompiler:
"""将 Transform 链编译为 Flink DataStream 操作"""
def compile(
self, transforms: list[TransformOperator],
source_stream: DataStream
) -> DataStream:
current = source_stream
for transform in transforms:
flink_op = transform.compile_flink()
if isinstance(flink_op, FlinkFilterFunction):
current = current.filter(flink_op)
elif isinstance(flink_op, FlinkMapFunction):
current = current.map(
flink_op,
output_type=self._to_flink_type(transform.output_schema())
)
elif isinstance(flink_op, FlinkKeyedAggregation):
current = (
current
.key_by(flink_op.key_selector)
.window(flink_op.window_assigner)
.aggregate(flink_op.aggregate_function)
)
elif isinstance(flink_op, FlinkJoin):
right_stream = self._resolve_stream(flink_op.right_ref)
current = (
current
.join(right_stream)
.where(flink_op.left_key)
.equal_to(flink_op.right_key)
.window(flink_op.window)
.apply(flink_op.join_function)
)
return current
#3.2 Flink 特有优化
Flink 特有优化:
1. 状态后端选择
- RocksDB StateBackend → 大状态(>1GB)
- HashMapStateBackend → 小状态(<1GB)
2. Checkpoint 配置
- interval: 1 分钟
- mode: EXACTLY_ONCE
- storage: MinIO (S3 兼容)
3. 水印策略
- BoundedOutOfOrdernessWatermarks
- maxOutOfOrderness: 5 seconds
4. 资源配置
- TaskManager memory: 4GB
- Parallelism: 根据 Kafka partition 数自动设置
#4. Spark 适配器
#4.1 DataFrame 编译
class SparkTransformCompiler:
"""将 Transform 链编译为 Spark DataFrame 操作"""
def compile(
self, transforms: list[TransformOperator],
source_df: DataFrame
) -> DataFrame:
current = source_df
for transform in transforms:
spark_op = transform.compile_spark()
if isinstance(spark_op, Column):
current = current.filter(spark_op)
elif isinstance(spark_op, SparkUDFMap):
udf = F.udf(spark_op.function, spark_op.return_type)
current = current.select(
udf(F.struct(*current.columns)).alias("result")
).select("result.*")
elif isinstance(spark_op, SparkGroupByAgg):
agg_exprs = [
getattr(F, func)(col).alias(alias)
for alias, (col, func) in spark_op.agg_exprs.items()
]
current = current.groupBy(
*spark_op.group_cols
).agg(*agg_exprs)
return current
#5. DuckDB 适配器
#5.1 SQL 链编译
class DuckDBTransformCompiler:
"""将 Transform 链编译为 DuckDB SQL"""
def compile(
self, transforms: list[TransformOperator],
source_sql: str
) -> str:
current_sql = f"({source_sql}) AS source"
for i, transform in enumerate(transforms):
sql_fragment = transform.compile_duckdb()
alias = f"t{i}"
if sql_fragment.startswith("WHERE"):
current_sql = (
f"(SELECT * FROM {current_sql} "
f"{sql_fragment}) AS {alias}"
)
elif sql_fragment.startswith("SELECT"):
current_sql = (
f"({sql_fragment} FROM {current_sql}) AS {alias}"
)
else:
current_sql = (
f"(SELECT * FROM {current_sql} "
f"{sql_fragment}) AS {alias}"
)
return f"SELECT * FROM {current_sql}"
#6. 引擎自动选择
class EngineSelector:
"""Transform 执行引擎自动选择"""
def select(self, pipeline: Pipeline) -> Engine:
features = self._analyze_pipeline(pipeline)
# 规则 1:CDC 源 → Flink
if features.has_cdc_source:
return Engine.FLINK
# 规则 2:窗口聚合 → Flink
if features.has_window_aggregation:
return Engine.FLINK
# 规则 3:数据量 < 100MB → DuckDB
if features.estimated_data_size_mb < 100:
return Engine.DUCKDB
# 规则 4:数据量 > 10GB → Spark
if features.estimated_data_size_mb > 10_000:
return Engine.SPARK
# 规则 5:Lambda 函数 + 大数据 → Flink
if features.has_lambda_transforms and features.estimated_data_size_mb > 100:
return Engine.FLINK
# 默认 → DuckDB(最轻量)
return Engine.DUCKDB
#7. 故障切换
class EngineFailover:
"""引擎故障切换"""
async def execute_with_failover(
self, pipeline: Pipeline, primary_engine: Engine
) -> ExecutionResult:
fallback_order = self._get_fallback_order(primary_engine)
for engine in [primary_engine] + fallback_order:
try:
result = await self._execute_on(pipeline, engine)
return result
except EngineUnavailableError as e:
logger.warning(
f"Engine {engine} unavailable: {e}. "
f"Trying next engine..."
)
except CompilationError as e:
logger.warning(
f"Cannot compile for {engine}: {e}. "
f"Trying next engine..."
)
raise AllEnginesFailedError(
f"All engines failed for pipeline '{pipeline.name}'"
)
def _get_fallback_order(self, primary: Engine) -> list[Engine]:
fallback_map = {
Engine.FLINK: [Engine.SPARK, Engine.DUCKDB],
Engine.SPARK: [Engine.FLINK, Engine.DUCKDB],
Engine.DUCKDB: [Engine.FLINK, Engine.SPARK],
}
return fallback_map[primary]
#8. 性能对比
三引擎性能基准对比:
测试场景:100 万行数据,filter → map → aggregate
┌──────────────────┬──────────┬──────────┬──────────┐
│ 指标 │ Flink │ Spark │ DuckDB │
├──────────────────┼──────────┼──────────┼──────────┤
│ 启动时间 │ 15s │ 30s │ 0.1s │
│ 处理时间 │ 8s │ 12s │ 3s │
│ 总时间(含启动) │ 23s │ 42s │ 3.1s │
│ 内存占用 │ 2GB │ 4GB │ 200MB │
│ 吞吐量(行/秒) │ 125K │ 83K │ 333K │
├──────────────────┼──────────┼──────────┼──────────┤
│ 流式延迟 (P99) │ 100ms │ 2s │ N/A │
│ 窗口聚合 │ 原生支持 │ 微批支持 │ 不支持 │
│ 容错恢复 │ 秒级 │ 分钟级 │ 无 │
└──────────────────┴──────────┴──────────┴──────────┘
结论:
- 批量小数据 → DuckDB 最优(快 7-14 倍)
- 流式处理 → Flink 最优(原生流引擎)
- 大规模批处理 → Spark 最优(分布式扩展)
#9. 测试策略
class TestTransformExecutor:
def test_filter_compiles_to_all_engines(self):
op = FilterOperator("age > 30")
assert op.compile_flink() is not None
assert op.compile_spark() is not None
assert "WHERE age > 30" in op.compile_duckdb()
def test_aggregate_consistency_across_engines(self):
"""三个引擎应产生相同结果"""
data = generate_test_data(1000)
op = AggregateOperator(
group_by=["type"],
aggregations={"total": ("value", "sum")}
)
flink_result = execute_flink(op, data)
spark_result = execute_spark(op, data)
duckdb_result = execute_duckdb(op, data)
assert_results_equal(flink_result, spark_result)
assert_results_equal(spark_result, duckdb_result)
def test_engine_auto_selection(self):
small_pipeline = Pipeline.create("small").source(csv_file("10mb.csv"))...
assert EngineSelector().select(small_pipeline) == Engine.DUCKDB
cdc_pipeline = Pipeline.create("cdc").source(mysql_cdc("db.table"))...
assert EngineSelector().select(cdc_pipeline) == Engine.FLINK
async def test_failover(self):
result = await failover.execute_with_failover(
pipeline, primary_engine=Engine.FLINK
)
assert result is not None # Should succeed on fallback engine
#Key Takeaways
-
统一 Transform 接口让用户与引擎解耦:同一个 filter/map/aggregate 定义可以编译到 Flink、Spark 或 DuckDB,用户无需关心底层引擎。
-
引擎自动选择基于数据特征:CDC 源 → Flink,小数据集 → DuckDB,大规模批处理 → Spark,系统根据数据量和处理需求自动决策。
-
DuckDB 在小数据场景下性能碾压:启动快(0.1s vs 15-30s)、处理快(单线程即可),是开发测试的最佳选择。
-
故障切换保障了生产可靠性:主引擎不可用时自动切换到备用引擎,避免因单一引擎故障导致 Pipeline 中断。
-
跨引擎结果一致性是核心约束:三个引擎对相同数据和 Transform 必须产生相同结果,这通过统一接口和一致性测试保证。
#Next Article
下一篇 S3-21《World Transform:全局数据一致性变换》 将介绍 Palantir Foundry 风格的 World Transform 概念——如何对整个 Ontology 世界状态进行原子性的全局变换。
Tags: #TransformExecutor #MultiEngine #Flink #Spark #DuckDB #EngineSelection #Failover #Consistency #智策平台 #coomia-dip #数据基座