首页
/ 使用 Pathway 与 Airbyte 构建流式 ETL 管道:从 GitHub 实时抽取、脱敏转换到落地存储

使用 Pathway 与 Airbyte 构建流式 ETL 管道:从 GitHub 实时抽取、脱敏转换到落地存储

2026-09-07 11:37:55作者:卓艾滢Kingsley

本指南以 docs/2.developers/7.templates/ETL/150.etl-python-airbyte.md 模板文档为主体,讲解如何在 Pathway Live Data Framework 中复用 Airbyte 的 350+ 开源源连接器完成实时 ETL:用 airbyte-serverless 配置抽取源、用 pw.io.airbyte.read 拉取增量数据流、用纯 Python 函数对流内 JSON 做 PII 脱敏转换,再通过 pw.io.jsonlines.write 等输出连接器加载到目标存储。读完本文你将掌握从 Airbyte 配置生成、流式读表参数调优,到基于 pw.apply 的无状态/有状态转换、以及多种落地连接器的完整端到端实战链路,并了解本地/远程(Google Cloud Run)两种连接器执行模式及其适用场景。

ETL 为什么需要流式化:从批处理到实时脱敏

在数据管理中,Extract(抽取)、Transform(转换)、Load(加载)三阶段决定了原始数据能否被有效利用:先从不同数据源抽取数据,再按业务需求与标准进行转换,最后加载到可分析、可支撑决策的目的地。传统 ETL 通常以批次为单位周期运行,数据延迟往往达到小时甚至天级。而当你既希望保留 ETL 的「先清洗、后入库」优势,又需要秒级延迟时,就需要流式 ETL:数据一旦产生,立即被抽取、即时转换、马上落地。

ETL(而非 ELT)特别适合不希望把原始数据直接存进数据仓库的场景。例如,包含个人身份信息(PII)的数据必须先在落地前被匿名化或脱敏,才能进入后续分析环节。本文以一个典型任务贯穿始终:以秒级延迟实时拉取指定 GitHub 仓库的提交(commits)流,然后在处理时移除 payload 中出现的所有邮箱。这里 Pathway 并非只能做过滤,通过连接、机器学习模型等更复杂的处理同样可行,脱敏只是演示流式转换最直观的入口。

架构上,本文的方案用两个开源工具分工:

  • Airbyte:开源数据集成平台,负责 ETL 中的 Extract 步骤,提供 350+ 输入连接器(GitHub、Postgres、Kafka、Stripe 等)供选用;
  • Pathway Live Data Framework:面向 Python 与 ML/AI 开发者的实时事件处理引擎,负责 Transform 与 Load,把 Airbyte 源产出的数据流变为可增量处理的实时表。

在仓库中,Pathway Airbyte 连接器用户指南 汇总了可用的免费档 Airbyte 源连接器清单;Airbyte 连接器实现连接器运行时逻辑 则是支撑本教程全部能力的真实源码。仓库还提供了一份可直接运行的同主题脚本 examples/projects/spark-data-preparation/main.py,演示了相同的 GitHub commits 流经脱敏后写入 Delta Lake(可扩展到 S3)的完整流程,可作为本文示例的落地版参照。

前置准备:安装 airbyte-serverless 与 Docker

本方案通过第三方工具 airbyte-serverless 管理 Airbyte 连接器的配置模板生成与本地执行,先通过 pip 安装:

pip install airbyte-serverless

注意:使用 airbyte-serverless 需要本机安装 Docker。Airbyte 连接器多以 Docker 镜像分发,本地执行时会拉起容器完成抽取。

Pathway 侧则直接以 import pathway as pw 方式接入,无需额外安装独立扩展——pw.io.airbyte.read 连接器是框架内置能力,其本地执行会优先尝试从 PyPI 安装对应 Python 连接器并放入独立虚拟环境,找不到 Python 包时才回退到 Docker 镜像执行(详见下文「连接器的两种执行方式」)。

第一步:用 airbyte-serverless 配置 GitHub 源

教程以 pathwaycom/pathway 仓库的 commits 流作为数据源。Airbyte 官方 GitHub 连接器 airbyte/source-github 即可读取该仓库。官方文档提供两种配置方式,任选其一。

方式 A:使用预定义配置模板

创建 YAML 文件 connections/github.yaml,将下列模板完整粘贴进去,并填入你自己的 token:

source:
  docker_image: "airbyte/source-github"  # 指定 Airbyte 连接器类型
  config:
    credentials:
      option_title: "PAT Credentials"  # 选用的认证方案名
      personal_access_token: <YOUR PERSONAL ACCESS TOKEN HERE>  # 取自 https://github.com/settings/tokens
    repositories:
      - pathwaycom/pathway  # Pathway 仓库
    api_url: "https://api.github.com/"
  streams: commits

关键字段说明:

  • docker_image:Airbyte 连接器的镜像标识,Pathway 据此推断连接器名称 github(源码中通过 source_config["docker_image"].removeprefix("airbyte/") 解析,见 python/pathway/io/airbyte/init.py),再尝试匹配 PyPI 包 airbyte-source-github
  • credentials.personal_access_token:GitHub Personal Access Token(PAT),需登录后前往 GitHub Tokens 页面生成,其 scope 需要 public_repo(访问公共仓库);
  • repositories:要读取的仓库列表,此处为 pathwaycom/pathway
  • streams:要订阅的 Airbyte 流,本示例固定为 commits

方式 B:手动生成配置模板(可选)

不同 Airbyte 源的参数、token 与 key 各不相同,若想基于连接器官方 spec 自动生成配置骨架,可用如下命令:

abs create github --source "airbyte/source-github"

执行后同样生成 ./connections/github.yaml,但此时里面是含大量注释与可选字段的完整 spec 模板,需要手工填值后再精简。官方文档强调:若不修改模板并移除未用字段,airbyte 会直接报错。手工精简步骤如下:

  1. 该连接器提供两套认证方式。选择较简单的 PAT 授权:先删除模板中处于未注释状态的 option_titleaccess_tokenclient_idclient_secret 字段,再取消注释名为 "Another valid structure for credentials" 的段落;该段落只要求一个带 public_repo scope 的 PAT;
  2. repositories 字段填入 pathwaycom/pathway
  3. 删除其它未使用的可选字段,即可运行。更完整的参数说明可参考 Airbyte GitHub 官方文档。

提示:仓库内 integration_tests/airbyte/test_airbyte.py 中的 _GITHUB_FAKE_CONFIG 给出了 GitHub 源配置的最小可解析形态(source.docker_imagesource.configstreams 之外只需要 repositoriesstart_date),可作为手动精简时的校验参考。

第二步:用 pw.io.airbyte.read 抽取数据流

配置就绪后即可写管道代码。首先导入 Pathway:

import pathway as pw

数据源的配置方式与其它输入连接器保持一致——调用 pw.io.airbyte.read,传入配置文件路径与要读取的流名列表:

commits_table = pw.io.airbyte.read(
    "./connections/github.yaml",
    streams=["commits"],
)

核心参数仅两个:config_file_path(指向上一节生成的 YAML)与 streams(要消费的 Airbyte 流名称)。GitHub 连接器把提交数据放在名为 commits 的流中,因此如上指定。

该读操作返回一张以 data 为列的 Pathway 表(实际由 _AirbyteRecordSchemadata: dict)承载,见 python/pathway/io/airbyte/init.py),列内是符合 Airbyte 格式的 JSON payload,可用 pw.Json 类型访问。

streaming 与 static 两种读取模式

上述代码会无限期运行:每当仓库产生新提交,数据会被追加进 commits_table。若只想读取当前已有的全部提交、处理完即正常退出(不等待新数据),把 mode 设为 "static" 即可:

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

mode 的两种取值语义(见 read 的 docstring):

  • "streaming"(默认):引擎按 refresh_interval 周期轮询新数据,持续追加;
  • "static":只消费当下可得的数据,一次性在一个 commit 中摄取完毕即终止。

refresh_interval 轮询频率控制

流式模式下,轮询源的频率由关键字参数 refresh_interval 控制,接受秒数datetime.timedelta默认 60 秒(源码 refresh_interval: DurationLike = 60)。注意:该参数在旧版本曾以 refresh_interval_ms 传递毫秒,现已被 refresh_interval 取代,若仍传入 refresh_interval_ms 会抛出提示改用秒或 timedelta 的 TypeError(源码第 337-356 行对此做了显式兼容拦截)。

对于 "remote" 执行方式(见后文),Google Cloud Run 按 vCPU 秒与 GiB 秒计费,过小的 refresh_interval 会显著放大运行次数与账单,远程模式下不建议设置过小值。

连接器的两种执行方式:local 与 remote

pw.io.airbyte.readexecution_type 参数决定连接器进程如何运行(默认 "local"),这一选择对部署形态与安全有直接影响:

  • "local"(默认):在本地执行连接器。Pathway 会先探测 PyPI 上是否存在同名 Python 包(如 airbyte-source-github),存在则以 venv 方式安装到隔离虚拟环境后运行;否则回退为 Docker 镜像执行。这一自动探测依赖 PyPI 可用性,官方强烈建议生产部署显式使用 enforce_method 固定行为,避免 PyPI 服务不可用时探测失败;
  • "remote":把连接器以 Google Cloud Run 任务方式跑在云端,需要额外提供 service_user_credentials_file(服务账号 JSON)与可选 gcp_region(默认 europe-west1)、gcp_job_name。该模式规避了容器嵌套调用 Docker 的问题并节省本机 CPU/内存。

关于 Docker-in-Docker 的工程警示

本地模式下,部分连接器必须用 Docker 镜像运行(如未提供 Python 包实现的连接器),这会导致"在 Docker 容器里再跑 Docker"(DinD)的经典难题。连接器的源码注释明确不推荐这种用法:DinD 需要 --privileged 启动特权容器,等于移除隔离边界、放大容器逃逸风险,同时嵌套 cgroup 会破坏 CPU/内存限额的准确统计,引发难以定位的 OOM。官方建议的替代方案是 enforce_method="pypi"(或默认的 venv 安装路径)将连接器当作 Python 包运行,或在资源隔离有保障的托管环境(如 GCP)中执行。

与之配套,enforce_method 取值 "docker"(强制走镜像)或 "pypi"(强制走 PyPI 包);dependency_overrides 可传入 pip requirement 字符串列表(如 ["airbyte-cdk==2.16.0"]),用于在连接器自带元数据未收紧传递依赖时锁定已知可用版本,仅对 PyPI 执行方式生效。仓库测试 integration_tests/airbyte/test_airbyte.py 恰好覆盖了这类约束:同一 GitHub 配置下,锁定坏版本 airbyte-cdk==7.13.0 会在导入期崩溃,而锁定 7.10.0 则能正常加载并命中 GitHub 的 401 鉴权错误——演示了 dependency_overrides 的典型排查用法。

实际代码调用示例(同时开启静态模式与 PyPI 强制):

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

仓库中的 examples/projects/spark-data-preparation/main.py 即采用上述写法读取 commits 后,再用 pw.apply 抽取 author.logincreated_at,最终写入 Delta Lake。

第三步:用纯 Python 转换实现实时流内脱敏

拿到 commits 流后,数据以 Pathway 表形式存在。接下来的任务是把 payload 中的邮箱全部移除。由于 GitHub 连接器返回的是 JSON,教程采用一个简单而有效的深度优先递归遍历算法:对 JSON 做 DFS,只要在非空白字符组成的连续串中出现 @ 字符,就整段剔除该非空串。这样能在保持实现极简的同时,以最大召回率确保没有任何 PII 残留。

转换函数完全是普通 Python,不调用任何框架 API:

import json


def remove_emails_from_data(payload):
    if isinstance(payload, str):
        # 字符串场景:按空格切分,剔除含 '@' 的片段后再合并回去
        return " ".join([item for item in payload.split(" ") if "@" not in item])

    if isinstance(payload, list):
        # 列表:对每个元素递归去邮箱后重组
        result = []
        for item in payload:
            result.append(remove_emails_from_data(item))
        return result

    if isinstance(payload, dict):
        # 字典:对每个 value 递归去邮箱(key 中不含邮箱)
        result = {}
        for key, value in payload.items():
            value = remove_emails_from_data(value)
            result[key] = value
        return result

    # 其余为原始类型:布尔、浮点、整数或 null,均无可删除内容
    return payload

递归函数处理三类容器:

  • str:按空格切分,丢弃含 @ 的片段再合并;
  • list:逐元素递归并重建列表;
  • dict:对每个 value 递归(key 理论不会含邮箱,故不处理);
  • 其它基础类型(bool/float/int/null):原样返回。

要把这个纯 Python 函数应用到表的每一行,使用 Pathway 的 pw.apply(将普通 Python 函数包装为逐行 UDF)。由于 pw.io.airbyte.read 产出的原始 payload 是 pw.Json 类型,需要先用 json.loads(payload.as_str()) 解析为 Python dict 再递归处理,随后把结果类型标注为 pw.Json 以接回表列:

def remove_emails(raw_commit_data: pw.Json) -> pw.Json:
    # 先把 pw.Json 解析成 Python dict
    data = json.loads(raw_commit_data.as_str())

    # 再应用递归方法删除邮箱
    return remove_emails_from_data(data)


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

执行后 commits_table 中已不含邮箱等个人信息,转换在流上实时生效——每个新到达的 commit 在进入表的同时即被脱敏。除 select 投影外,examples/projects/spark-data-preparation/main.py 还示范了用 pw.apply 抽取嵌套字段(extract_author_loginextract_commit_timestamp),说明同样的 UDF 机制可扩展到字段拆分、类型化时间解析等更丰富的转换。

第四步:用输出连接器落地脱敏后的流

转换完成后即可把结果写出。Pathway 提供丰富的输出连接器可供选择,部分典型选项包括:

  • Kafka 中的 topic(pw.io.kafka.write);
  • Logstash 端点;
  • Postgres 表(pw.io.postgres.write);
  • 甚至是 Python 回调(pw.io.subscribe)。

最简单的本地磁盘方案是 jsonlines 输出连接器,把每条记录以 JSON Lines 形式写入本地文件:

pw.io.jsonlines.write(commits_table, "commits.jsonlines")

上述代码会把结果持续写入 commits.jsonlines。在流式(streaming)模式下,文件会随着新 commit 到达而持续增长;若改用 mode="static",则写入既定数据后即可结束。

最后,别忘了启动引擎:

pw.run()

至此,数据已在流上实时完成匿名化并落盘。这是一个最小但完整的 ETL 闭环:airbyte-serverless 配置抽取 → pw.io.airbyte.read 拉流 → 纯 Python 递归脱敏 → pw.io.jsonlines.write 落地。真实场景中,把输出的 jsonlines 换成仓库 examples/projects/spark-data-preparation/main.py 演示的 pw.io.deltalake.write(本地或经 pw.io.s3.AwsS3Settings 写 S3)即可对接数据湖。

从示例到生产:可调参数与注意事项速查

下表汇总 pw.io.airbyte.read 在生产中常用的参数(依据 python/pathway/io/airbyte/init.py 的签名与注释整理):

参数 默认值 作用与建议
config_file_path 必填 airbyte-serverless(或 Pathway CLI)生成的 YAML 路径,其 source 段须预先配好
streams 必填 要读取的 Airbyte 流名列表
execution_type "local" 本地执行,或 "remote"(Google Cloud Run 任务)
mode "streaming" "streaming" 按周期轮询持续追加;"static" 一次性读尽现有数据后终止
refresh_interval 60 秒 流式模式轮询间隔,秒数或 datetime.timedelta
enforce_method 自动探测 "docker" 强制走镜像、"pypi" 强制走 Python 包;生产建议显式指定
dependency_overrides 以 pip requirement 钉住连接器虚拟环境内的传递依赖版本(仅 PyPI 方式生效)
service_user_credentials_file "remote" 执行必需的 Google 服务账号 JSON
gcp_region / gcp_job_name europe-west1 / 自动 Cloud Run 任务的区域与任务名
name 自动 连接器在日志与监控面板中的唯一名;开启持久化时用作快照名
max_backlog_size 限制读入待处理条目的积压上限,防止大源首波爆发导致内存尖峰

几点工程提醒:

  1. 同步模式一致性:所有 streams 必须具有相同的 sync_modeincrementalfull_refresh),_construct_local_source 会在混用时报 ValueError;Pathway 侧对 Airbyte 协议的状态记录(LEGACY / GLOBAL / STREAM / PER_STREAM 四种格式)均有解析支持,断点续读的状态通过 _airbyte_states 记录上报(见 python/pathway/io/airbyte/logic.py)。
  2. 持久化限制:启用持久化时,当前 Airbyte 连接器仅支持单流场景,多流需要拆成多个连接器、每个连接器只读一个流。
  3. 失败重试:抽取失败时连接器按指数退避自动重试,最多 5 次后抛出异常(源码中的 MAX_RETRIES = 5)。
  4. 认证安全:GitHub PAT 等敏感凭据直接存在于 YAML 中,仓库已有 integration_tests/airbyte/.gitignore 相关约定 之外的版本控制场景请自行做好密钥管理(环境变量注入可用 env_vars 参数传入连接器)。

小结:350+ 数据源即刻接入实时管道

本教程以一个「实时拉取 GitHub commits 并即时脱敏」的最小用例走通了流式 ETL 的完整链路:airbyte-serverless 承担源配置与连接器编排,Pathway 负责增量拉流、实时转换与多目标落地。由于 Pathway 把 Airbyte 的每个流抽象成一张可持续追加的实时表,并允许在 Python 中任意组合 selectpw.apply 乃至连接操作,其可扩展空间远不止脱敏——Airbyte 免费档的 350+ 源连接器(数据库、SaaS、对象存储等)理论上均可通过同一套 pw.io.airbyte.read 接入,再搭配 JSONLines、Delta Lake、Postgres、Kafka 等输出端,即可把「抽取-转换-加载」从批处理范式平滑迁移到秒级实时管道。

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

项目优选

收起
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++
916
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