使用 Pathway 与 Airbyte 构建流式 ETL 管道:从 GitHub 实时抽取、脱敏转换到落地存储
本指南以 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 会直接报错。手工精简步骤如下:
- 该连接器提供两套认证方式。选择较简单的 PAT 授权:先删除模板中处于未注释状态的
option_title、access_token、client_id、client_secret字段,再取消注释名为 "Another valid structure for credentials" 的段落;该段落只要求一个带public_reposcope 的 PAT; - 在
repositories字段填入pathwaycom/pathway; - 删除其它未使用的可选字段,即可运行。更完整的参数说明可参考 Airbyte GitHub 官方文档。
提示:仓库内 integration_tests/airbyte/test_airbyte.py 中的
_GITHUB_FAKE_CONFIG给出了 GitHub 源配置的最小可解析形态(source.docker_image、source.config与streams之外只需要repositories与start_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 表(实际由 _AirbyteRecordSchema(data: 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.read 的 execution_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.login 与 created_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_login、extract_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 |
无 | 限制读入待处理条目的积压上限,防止大源首波爆发导致内存尖峰 |
几点工程提醒:
- 同步模式一致性:所有
streams必须具有相同的sync_mode(incremental或full_refresh),_construct_local_source会在混用时报ValueError;Pathway 侧对 Airbyte 协议的状态记录(LEGACY / GLOBAL / STREAM / PER_STREAM 四种格式)均有解析支持,断点续读的状态通过_airbyte_states记录上报(见 python/pathway/io/airbyte/logic.py)。 - 持久化限制:启用持久化时,当前 Airbyte 连接器仅支持单流场景,多流需要拆成多个连接器、每个连接器只读一个流。
- 失败重试:抽取失败时连接器按指数退避自动重试,最多 5 次后抛出异常(源码中的
MAX_RETRIES = 5)。 - 认证安全:GitHub PAT 等敏感凭据直接存在于 YAML 中,仓库已有 integration_tests/airbyte/.gitignore 相关约定 之外的版本控制场景请自行做好密钥管理(环境变量注入可用
env_vars参数传入连接器)。
小结:350+ 数据源即刻接入实时管道
本教程以一个「实时拉取 GitHub commits 并即时脱敏」的最小用例走通了流式 ETL 的完整链路:airbyte-serverless 承担源配置与连接器编排,Pathway 负责增量拉流、实时转换与多目标落地。由于 Pathway 把 Airbyte 的每个流抽象成一张可持续追加的实时表,并允许在 Python 中任意组合 select、pw.apply 乃至连接操作,其可扩展空间远不止脱敏——Airbyte 免费档的 350+ 源连接器(数据库、SaaS、对象存储等)理论上均可通过同一套 pw.io.airbyte.read 接入,再搭配 JSONLines、Delta Lake、Postgres、Kafka 等输出端,即可把「抽取-转换-加载」从批处理范式平滑迁移到秒级实时管道。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0629
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00