首页
/ Pathway + Delta Lake 构建面向 Spark 分析的 ETL 数据准备管道

Pathway + Delta Lake 构建面向 Spark 分析的 ETL 数据准备管道

2026-09-07 14:39:14作者:庞队千Virginia

本篇技术指南以仓库模板文档 175.delta_lake_etl.md 为核心,系统讲解如何使用 Pathway(Pathway Live Data Framework)从 GitHub 拉取提交历史,经过脱敏与列结构整理后,以 Delta Lake 格式落盘到本地文件系统或 S3,并最终交给 Spark 做分析查询。读完本文,你将掌握 pw.io.airbyte.read 数据抽取、pw.apply 驱动的 Python UDF 数据清洗与特征抽取、pw.io.deltalake.write 结果落盘以及 Delta Lake + Spark 会话配置与查询的完整实战链路,可直接复刻到自己的数据分析业务中。

为什么用 Pathway + Delta Lake 为 Spark 准备数据

Apache Spark 是企业中应用极为广泛的分布式数据分析引擎,而 Delta Lake 为 Spark 补齐了数据可靠性短板:它基于 Parquet 列式存储之上提供 ACID 事务、可扩展的元数据管理和 time travel(数据版本回放) 能力,用户随时可以访问和查询表的任意历史版本。

不过 Spark 本身并不擅长"把分散的原始数据搬进来并整理好"这件事。Pathway 在这里扮演的是 数据接入与预处理的 ETL 层:从几百种数据源完成抽取(Extraction)与转换(Transformation),再以 Delta Lake 格式写入本地文件系统或 S3,让 Spark 分析时面对的是整洁、安全、结构化的表。仓库中 spark-data-preparation 示例工程即本文讲解的完整可运行代码,由 main.py 与配套配置、依赖清单和 Dockerfile 组成。

从源码看,Delta Lake 输出连接器位于 python/pathway/io/deltalake/init.py,其 pw.io.deltalake.write(见 该文件 #L527-L540)会依据 URI 自动判断存储类型:以 s3://s3a:// 开头的路径走 S3,其余路径使用本地文件系统。

示例任务与三步 ETL 架构

教程要构建的管道非常贴近真实的"隐私合规 + 分析就绪"场景:

  • 数据来源:GitHub 仓库的提交(commit)历史流;
  • 目标:把提交数据整理成可供 Spark 分析的表格;
  • 约束:在分析前必须移除提交中可能携带的邮箱等个人身份信息(PII),保护贡献者隐私。

管道遵循经典 ETL 结构,拆解为三步:

  1. 抽取(Extraction):借助基于 Airbyte Serverless 的 Airbyte 连接器配置数据摄取,从 GitHub 仓库拉取提交数据;
  2. 转换(Transformation):编写 Python 用户自定义函数(UDF),剔除邮箱这类个人数据,并按分析需要整理/新增数据列;
  3. 加载(Loading):使用 Pathway 的 Delta Lake 连接器把整理后的数据写入 Delta Lake。

下面按这三步逐一展开,最后再用 Spark 对产物做查询验证。

第一步 · 抽取:用 Airbyte GitHub 连接器接入提交流

生成连接器配置

抽取的第一步是准备 Airbyte 的 source 配置。Pathway 内置的 Airbyte 封装基于 airbyte-serverless,可直接用 CLI 生成某个连接器的样例配置:

pathway airbyte create-source github --image airbyte/source-github:latest

生成的配置中只有少数几个字段需要实际编辑:

字段 含义与填法
personal_access_token GitHub PAT(Personal Access Token),连接 GitHub 最直接的方式,在 GitHub 的 Tokens 设置页生成后填入
repositories 要做分析的仓库列表,示例中填入 pathwaycom/pathway
streams 本次目标是分析提交流,因此填写 commits

配置编辑完成后形如下方 YAML。注意:示例工程会把该配置保存在 ./github-config.yaml 中,管道运行时读取的就是这个文件(真实文件见 github-config.yaml)。

source:
  docker_image: "airbyte/source-github:latest"
  config:
    credentials:
      option_title: "PAT Credentials"
      personal_access_token: YOUR_PAT_TOKEN
    repositories:
      - pathwaycom/pathway
    api_url: "https://api.github.com/"
  streams: commits

pw.io.airbyte.read 建立输入表

有了配置文件,下一步是调用 Pathway 的 Airbyte 连接器 pw.io.airbyte.read 把数据接入为一张表。它的完整签名定义在 python/pathway/io/airbyte/init.py 中(见 read 函数签名 #L112-L127)。本任务需要指定两个关键参数:

  1. Config File:上一步生成的 ./github-config.yaml。原文此处写作 ./github-commits.yaml 属笔误,示例工程中实际文件名与读取路径均为 github-config.yaml
  2. Mode:测试时使用 static 模式,即程序在启动后一次性读完全部现存提交便结束运行,便于本地验证。

此外,GitHub 连接器本身是 Python 实现,因此无需 Docker即可运行:通过参数 enforce_method="pypi" 让 Pathway 直接从 PyPI 安装并使用 airbyte-source-github 库,而不是拉取 Docker 镜像。

得到如下代码:

import pathway as pw

commits_table = pw.io.airbyte.read(
    "./github-config.yaml",
    streams=["commits"],
    enforce_method="pypi",
    mode="static",
)

源码级补充:连接器执行语义。 从 read 的 docstring 可以看到更多工程化细节(见 python/pathway/io/airbyte/init.py#L129-L195):

  • 默认情况下 execution_type="local",Pathway 会先尝试在 PyPI 找到对应连接器并在独立虚拟环境中安装最新版;找不到时回退到 Docker 镜像执行。正因为存在这种自动探测,官方强烈建议在生产环境中显式指定 enforce_method,避免 PyPI 服务不可用时探测失败;
  • mode 支持 "streaming""static" 两种取值:"streaming"(默认)每隔 refresh_interval 秒(默认 60 秒)轮询一次新数据;"static" 则只把当前可用数据一次性以单个 commit 全部摄取。本文演示在 static 下运行,管道跑完即终止;
  • 返回的表固定带有一个 data 列,其类型为 pw.Json,内容格式与 Airbyte 产出的原始 JSON 一致——这正是下一步转换逻辑的输入;
  • 生产环境若身处容器内(Docker-in-Docker),docstring 亦明确提醒不应通过 Docker 方式运行连接器,这同样解释了本例选用 pypi 方法的另一个动机;
  • 可选参数还包括 execution_type="remote"(将连接器作为 Google Cloud Run 任务运行)、name(用于日志与监控标识,开启持久化时作为快照名)、max_backlog_size(限制同时驻留内存的输入条目数,避免大源首波突发数据撑爆内存)等。

第二步 · 转换:Python UDF 清洗 PII 并整理分析列

数据抽取完成后进入转换阶段。本任务的转换分两层目的:隐私合规(剔除邮箱)与分析就绪(抽出 Spark 日常分析要用的关键列)。

UDF 一:递归清洗 JSON 中的邮箱

核心思路是对整条 commit 的 JSON 递归扫描,凡字符串字段里包含 @(邮箱特征)即视为 PII 并剔除,从而让数据可以在组织内安全地分发给各数据分析师。该方法的完整实现如下:

import json


def remove_emails_from_data(payload):
    if isinstance(payload, str):
        # The string case is obvious: it's getting split and then merged back after
        # the email-like substrings are removed
        return " ".join([item for item in payload.split(" ") if "@" not in item])

    if isinstance(payload, list):
        # If the payload is a list, one needs to remove emails from each of its
        # elements and then return the result of the processing
        result = []
        for item in payload:
            result.append(remove_emails_from_data(item))
        return result

    if isinstance(payload, dict):
        # If the payload is a dict, one needs to remove emails from its keys and
        # values and then return the clean dict
        result = {}
        for key, value in payload.items():
            # There are no e-mails in the keys of the returned dict
            # So, we only need to remove them from values
            value = remove_emails_from_data(value)
            result[key] = value
        return result

    # If the payload is neither str nor list or dict, it's a primitive type:
    # namely, a boolean, a float, or an int. It can also be just null.
    #
    # But in any case, there is no data to remove from such an element.
    return payload


def remove_emails(raw_commit_data: pw.Json) -> pw.Json:
    # First, parse pw.Json type into a Python dict
    data = json.loads(raw_commit_data.as_str())

    # Next, just apply the recursive method to delete e-mails
    return remove_emails_from_data(data)

这里有一个值得注意的 Pathway 类型细节:输入是 pw.Json,无法直接当普通 dict 遍历,因此先用 json.loads(raw_commit_data.as_str()) 把它转成 Python dict;对 dict 只需清洗 value 即可,因为返回 dict 的 key 中不会出现邮箱。

清洗后的数据需要重新写回表,使用 pw.apply 把普通 Python 函数提升为可在 Pathway 数据流上执行的 UDF:

commits_table = commits_table.select(
    data=pw.apply(remove_emails, pw.this.data)
)

UDF 二:抽取 author_logincommit_timestamp

原始 commit JSON 是嵌套结构,Spark 做"按贡献者统计活跃度"这类每日分析时需要扁平化的列。定义两个 UDF 分别抽取贡献者登录名与变更时间戳:

def extract_author_login(commit_data: pw.Json) -> str:
    if not commit_data["author"]:
        return ""
    return commit_data["author"]["login"]


def extract_commit_timestamp(commit_data: pw.Json) -> pw.DateTimeUtc:
    return pw.DateTimeUtc(commit_data["created_at"].as_str())


commits_table = commits_table.select(
    author_login=pw.apply(extract_author_login, pw.this.data),
    commit_timestamp=pw.apply(extract_commit_timestamp, pw.this.data),
    data=pw.this.data,
)

代码要点:

  • author 字段为空(如某些系统机器人提交)时,extract_author_login 返回空字符串兜底,避免后续统计出现空值中断;
  • 时间字符串通过构造器 pw.DateTimeUtc(...) 转换为 Pathway 的统一时间类型,保证后续时间窗口分析的类型正确性;在示例工程 main.py 中,登录名抽取还显式调用了 .as_str() 完成 pw.Jsonstr 的取值转换;
  • select 阶段同时保留原始 data 列——Spark 侧若需要深挖原始信息,仍然可以访问未被破坏的完整 JSON。

源码级补充:UDF 的运行方式。 使用 pw.apply 而非在 Python 里写循环,是因为 Pathway 把数据处理建模为持续更新的数据流:UDF 会以引擎原生的方式被调度执行,即使将来把 modestatic 切回 streaming、让数据持续涌入,同一套转换逻辑也无需改写即可增量地处理每个新提交,这正是"数据准备层"面向实时数据源的可扩展性所在。

第三步 · 加载:写入 Delta Lake

数据准备完毕,进入落盘环节。重要前提:Pathway 的 Delta Lake 写入连接器仅在 Pathway Live Data Framework 的 ScaleEnterprise 商业层级中提供,需要先从 Pathway 官网申请许可证密钥,再通过 pw.set_license_key 注入运行时:

pw.set_license_key("YOUR_LICENSE_KEY")

这与源码中 pw.io.deltalake.write 第一行即执行 _check_entitlements("deltalake")(见 python/pathway/io/deltalake/init.py#L647)的授权检查逻辑相互印证。若在容器中运行示例工程,许可证也可直接以环境变量形式注入:Dockerfile 的启动方式要求运行时提供 PATHWAY_LICENSE_KEY(见 Dockerfile示例 README)。

场景一:写入本地文件系统

本地路径最省事:Delta 表的 schema 会根据 Pathway 表的列自动推断,因此只需要给出落盘目录即可。

pw.io.deltalake.write(commits_table, "./commit-storage")

从源码看,write 的可选参数远比表面丰富(见 deltalake 连接器 #L554-L593),理解它们有助于在生产中调优:

  • partition_columns:建表时的分区列,按需做分区裁剪加速 Spark 查询;
  • min_commit_frequency:两次存储提交之间的最小时间间隔(毫秒),默认 60_000。设 None 时每个 finalize 的 mini-batch 都会尽快提交。注意:Delta Lake 每次 commit 都会新产生一个文件并在事务日志写一条记录,频繁提交会增加产物表后续处理的开销,因此文档建议限制提交频率,并配合 Delta 的 vacuum / optimize 操作进一步合并文件、控制表内 chunk 数量;
  • sort_by:在每个 mini-batch 内按指定列做升序排序,多列时按值元组字典序比较,可让相邻行尽量落在一起、提升压缩与扫描效率;
  • output_table_type:默认 "stream_of_changes" 输出变更流(额外带 timediff 两列,分别表示计算批号与增删标记);设为 "snapshot" 则以每个 mini-batch 原子更新的方式维护当前全量状态(会附加 _id 字段,删除频繁时性能下降,且不适合超出内存的表);
  • table_optimizer:输出表的优化规则参数。

另外,Pathway 建表时会把自身的 schema 写入列元数据,因此这张表之后可以被 pw.io.deltalake.read 免显式 schema 地直接读回,保证"写入—回读"闭环顺畅(见 deltalake 连接器 #L549-L552)。

场景二:写入 S3 Bucket

S3 落盘复杂在凭证管理,官方区分两种主场景:

  1. 已认证的 AWS 机器:若运行管道的主机已完成 AWS 认证,直接给一个 s3:// 路径即可,无需显式凭证——源码同样支持该简化路径(writes3_connection_settings 缺省但 URI 为 S3 前缀时会尝试从路径推导并授权,见 python/pathway/io/deltalake/init.py#L47-L55);
  2. 云部署 / 无默认凭证环境:需要借助 pw.io.s3.AwsS3Settings 显式提供桶凭证,并建议把密钥放环境变量而非硬编码:
import os


# Forming the credentials
# To protect credentials, it's advised to store
# access key and secret access key in the environment variables
credentials = pw.io.s3.AwsS3Settings(
    access_key=os.environ["AWS_S3_ACCESS_KEY"],
    secret_access_key=os.environ["AWS_S3_SECRET_ACCESS_KEY"],
    bucket_name="aws-integrationtest",
    region="eu-central-1",
)

# Using the credentials to write the Delta Table
delta_table_path = "s3://bucket-name/your/delta/table/path/"
pw.io.deltalake.write(
    commits_table,
    delta_table_path,
    s3_connection_settings=credentials,
)

示例工程对 S3 场景的完整封装在 main.py:以 AWS_S3_OUTPUT_PATH 环境变量作为可选开关,存在时才构建 AwsS3Settings 并执行第二次写入;对应的容器启动命令需要一次性注入 PATHWAY_LICENSE_KEYAWS_S3_OUTPUT_PATHAWS_S3_ACCESS_KEYAWS_S3_SECRET_ACCESS_KEYAWS_BUCKET_NAMEAWS_REGION 六个变量(完整命令见 示例 README)。注意 AwsS3Settings 除了密钥与 region,还支持自定义 endpoint——托管在 AWS 之外、兼容 S3 协议的对象存储(如 MinIO 等)也需要通过它来指定服务地址。

运行管道并检查产物结构

一切就绪后,用 pw.run 启动计算。示例中关闭了监控输出以保持日志干净:

pw.run(monitoring_level=pw.MonitoringLevel.NONE)

运行结束后,在 UNIX 终端用 find 查看本地生成的 lake 目录结构:

find ./commit-storage

输出如下:

./commit-storage
./commit-storage/part-00000-06ff5709-5390-47a4-a5b5-72b3b0fa1970-c000.snappy.parquet
./commit-storage/_delta_log
./commit-storage/_delta_log/00000000000000000000.json
./commit-storage/_delta_log/00000000000000000001.json

两个组成部分各自的作用:

  • part- 前缀的 Parquet 文件:数据块被转换成的列式存储文件。Parquet 之所以关键,在于其针对查询场景优化的读性能——本质上是一批小型只读数据库,允许快速高效的数据检索,是管道效率的重要来源;
  • _delta_log 事务日志目录:存放 lake 的事务日志,支撑 Delta Lake 的 ACID 事务、可扩展元数据管理与数据版本化。本教程中日志包含两个以版本号命名的 JSON 文件:
    • 版本 0:记录建表细节,包括表 schema;
    • 版本 1:记录 Parquet 数据块被追加进表。

S3 场景同样可以用 AWS 命令行工具查看 Delta Table 内容:

aws s3 ls s3://bucket-name/your/delta/table/path/

返回:

                           PRE _delta_log/
2024-07-19 13:36:18     325161 part-00000-6cd0a8b6-2fb9-494e-acef-2856988bb20c-c000.snappy.parquet

与本地情形一致:顶层是一个 Parquet 文件加一个 _delta_log 目录。想看事务日志,把路径后追加 _delta_log/ 即可:

aws s3 ls s3://bucket-name/your/delta/table/path/_delta_log/

返回:

2024-07-19 13:36:11       1843 00000000000000000000.json
2024-07-19 13:36:18       5841 00000000000000000001.json

同样是两个版本文件:一个记录建表、一个记录数据追加。这套"版本化日志 + 数据文件"的布局正是 Spark 上做 ACID 写入和 time travel 查询的基础设施。

用 Spark 查询产物,闭环验证

管道产出验证的最后一步是回到 Spark 侧:安装 Spark 相关库并实现基本查询,确认"Pathway 准备的数据能被 Spark 直接消费"。

安装依赖并配置 Spark 会话

需要两个库:pyspark(计算引擎核心库)与 delta-spark(让 Spark 认识 Delta Lake 的扩展)。示例工程的 requirements.txt 同时把二者与 pathway 一起声明:

pip install pyspark
pip install delta-spark

接着配置 Spark 会话——核心是注册 Delta 的 SQL 扩展与 Catalog 实现:

from delta import configure_spark_with_delta_pip

from pyspark.sql import SparkSession
from pyspark.sql.functions import countDistinct

builder = (
    SparkSession.builder.appName("DeltaLakeUserLoginsAnalysis")
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension")
    .config(
        "spark.sql.catalog.spark_catalog",
        "org.apache.spark.sql.delta.catalog.DeltaCatalog",
    )
)
spark = configure_spark_with_delta_pip(builder).getOrCreate()

configure_spark_with_delta_pip 会保证 Spark 运行时能拿到 Delta Lake 对应的 JAR 包依赖。

执行查询:统计仓库唯一贡献者

例如,从准备好的数据中提取贡献者登录名列表:

# Read the table
df = spark.read.format("delta").load("./commit-storage")

# Get only unique committer logins
unique_logins_df = df.select("author_login").distinct()

# Display the logins without length limit
unique_logins_df.show(unique_logins_df.count())

# Gracefully stop the Spark session
spark.stop()

执行后确实输出该仓库的全部唯一贡献者:

+------------------+
|      author_login|
+------------------+
|    "szymondudycz"|
|      "gitFoxCode"|
|   "pw-ppodhajski"|
|     "zxqfd555-pw"|
|"krzysiek-pathway"|
|         "embe-pw"|
|    "janchorowski"|
|         "cla-ire"|
|          "olruas"|
|          "izulin"|
| "dependabot[bot]"|
|       "Saksham65"|
|     "Pathway-Dev"|
|         "XGendre"|
|         "avriiil"|
|                  |
|         "dxtrous"|
|        "ananyaem"|
|        "voodoo11"|
|        "lewymati"|
| "KamilPiechowiak"|
|   "berkecanrizai"|
+------------------+

注意输出中含一行空登录名——这正是前面 extract_author_login 对无 author 字段提交返回空字符串的兜底结果,在真实数据中常见,统计时可用 WHERE author_login <> '' 之类的条件过滤。需要对比验证时,可直接参考 main.py:示例工程把 Pathway 写入与 Spark 读取放在同一脚本里顺序执行,落盘完成后立即回读验证。

在仓库中复现本教程

若想一键复现本文全部流程,仓库已给出容器化版本:

  1. 获取 GitHub PAT 并填入 github-config.yamlpersonal_access_token 字段(文件内有注释指引);
  2. 构建镜像:docker build --no-cache -t spark-data-preparation .
  3. 运行容器,并通过 -e PATHWAY_LICENSE_KEY=YOUR_LICENSE_KEY 注入许可证(Delta Lake 连接器仅随 Scale/Enterprise 层级提供): docker run -e PATHWAY_LICENSE_KEY=YOUR_LICENSE_KEY -t spark-data-preparation

镜像构建时会安装默认 JRE 并 pip install pathwaypysparkdelta-spark 三个依赖(见 Dockerfile);S3 场景则在启动时追加 AWS_S3_OUTPUT_PATHAWS_S3_ACCESS_KEYAWS_S3_SECRET_ACCESS_KEYAWS_BUCKET_NAMEAWS_REGION 等环境变量即可(见 示例 README)。

小结

本教程完整展示了"用 Pathway 做 ETL 预处理、用 Delta Lake 做存储格式、交给 Spark 做分析"的分层协作范式:

  • 抽取:一条 pathway airbyte create-source 命令生成连接器配置,一次 pw.io.airbyte.read 调用即可接入 GitHub 提交流,Python 型连接器配合 enforce_method="pypi" 甚至可以摆脱 Docker 依赖;
  • 转换pw.Json + pw.apply 让普通 Python 函数成为流式 UDF,既能递归剔除邮箱这类 PII 满足隐私合规,也能抽出 author_logincommit_timestamp 等 Spark 高频使用的扁平列;
  • 加载pw.io.deltalake.write 一条语句同时支持本地路径与 S3(含兼容 S3 的自建对象存储),自动建表、自动推断 schema、自动维护事务日志,且写入层由授权检查保证商业能力边界;
  • 验证pyspark + delta-spark 组合只需少量会话配置即可对产物执行标准 Delta Lake 查询,证明数据已被整理成 Spark 可直接消费的分析就绪形态。

对任何以 Spark 为分析底座、又希望把数据接入与预处理从批处理脚本中解放出来的团队而言,这套"Pathway 负责搬与洗、Delta Lake 负责存与管、Spark 负责查与分析"的组合,是一条上手门槛低、数据安全边界清晰、可平滑扩展到持续数据流的实践路径。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.13 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
897
5.8 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
529
593
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
915
1.83 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.58 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.35 K
1.46 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.01 K
515
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
547
388