返回博客

Transform 执行器:多引擎适配层

Tags: #TransformExecutor #MultiEngine #Flink #Spark #DuckDB #智策平台

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

系列: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 不同引擎的优势场景

Code
三大引擎优势对比:

┌─────────────────┬──────────────┬──────────────┬──────────────┐
│ 场景              │ Flink        │ Spark        │ DuckDB       │
├─────────────────┼──────────────┼──────────────┼──────────────┤
│ 实时流处理        │ ★★★★★       │ ★★★          │ ★             │
│ CDC 数据摄入      │ ★★★★★       │ ★★           │ ★             │
│ 大规模批处理      │ ★★★          │ ★★★★★       │ ★★            │
│ 交互式分析        │ ★★           │ ★★★          │ ★★★★★        │
│ 小数据集处理      │ ★★           │ ★★           │ ★★★★★        │
│ 窗口聚合         │ ★★★★★       │ ★★★★         │ ★★            │
│ 嵌入式执行        │ ✗            │ ✗            │ ★★★★★        │
│ 资源占用         │ 高            │ 高            │ 极低           │
└─────────────────┴──────────────┴──────────────┴──────────────┘

#1.2 统一接口的价值

Code
统一 Transform 接口的价值:

用户视角:
  pipeline.transform(filter(lambda r: r["age"] > 30))
  → 不关心底层是 Flink/Spark/DuckDB

系统视角:
  → 根据数据量、延迟要求和资源状况自动选择引擎
  → 一个引擎故障时自动切换到另一个
  → 开发/测试用 DuckDB(秒级启动),生产用 Flink

#2. Transform 接口定义

#2.1 统一接口

Python
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.1 DataStream 编译

Python
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
Code
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 编译

Python
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 链编译

Python
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. 引擎自动选择

Python
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. 故障切换

Python
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. 性能对比

Code
三引擎性能基准对比:

测试场景: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. 测试策略

Python
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

  1. 统一 Transform 接口让用户与引擎解耦:同一个 filter/map/aggregate 定义可以编译到 Flink、Spark 或 DuckDB,用户无需关心底层引擎。

  2. 引擎自动选择基于数据特征:CDC 源 → Flink,小数据集 → DuckDB,大规模批处理 → Spark,系统根据数据量和处理需求自动决策。

  3. DuckDB 在小数据场景下性能碾压:启动快(0.1s vs 15-30s)、处理快(单线程即可),是开发测试的最佳选择。

  4. 故障切换保障了生产可靠性:主引擎不可用时自动切换到备用引擎,避免因单一引擎故障导致 Pipeline 中断。

  5. 跨引擎结果一致性是核心约束:三个引擎对相同数据和 Transform 必须产生相同结果,这通过统一接口和一致性测试保证。

#Next Article

下一篇 S3-21《World Transform:全局数据一致性变换》 将介绍 Palantir Foundry 风格的 World Transform 概念——如何对整个 Ontology 世界状态进行原子性的全局变换。

Tags: #TransformExecutor #MultiEngine #Flink #Spark #DuckDB #EngineSelection #Failover #Consistency #智策平台 #coomia-dip #数据基座