首页
/ Pathway 自定义 Python Connector 实战:用 `ConnectorSubject` 把 Twitter 等任意流式数据接入实时数据管道

Pathway 自定义 Python Connector 实战:用 `ConnectorSubject` 把 Twitter 等任意流式数据接入实时数据管道

2026-09-07 13:09:06作者:裴锟轩Denise

本文基于 Pathway 官方示例仓库中的 Twitter 自定义连接器(examples/projects/custom-python-connector-twitter)展开,讲解如何通过继承 pw.io.python.ConnectorSubject 编写自己的 Python 输入连接器,把外部数据源(文中以 Twitter 流为例)变成一张实时变化的 Pathway 表并持续写出到 CSV。读完本文,你将掌握自定义 Python 连接器的生命周期回调、数据投喂方法、Schema 声明与批量提交语义,并能把同一套模式复用到任意可被 Python 拉取的流式数据源上。

示例背景与适用场景

示例 README(examples/projects/custom-python-connector-twitter/README.md)开篇给出了一个重要前提:

⚠️ Twitter has turned off its free tier streaming API.(Twitter 已关闭其免费版流式 API。)

也就是说,Twitter 在这里只是一个演示用的数据源。README 明确指出,示例的核心目的是展示如何实现自定义 Python 连接器,Twitter API 仅作为例子——任何数据流都可以通过同样的方式送入 Pathway 的实时数据框架(Live Data Framework)。

因此,本示例的学习价值并不依赖 Twitter 是否可用,而在于它完整演示了"连接器(Subject)→ pw.io.python.read → Pathway 表 → 输出"这条自定义数据接入链路的全部写法。

整条管线的运行逻辑

示例脚本 twitter_connector_example.py 把一条"推文流"变成了实时表并落盘到 output.csv,整体结构如下:

Twitter 流 (tweepy.StreamingClient.sample)
        │  on_response 回调
        ▼
TwitterSubject 继承 pw.io.python.ConnectorSubject
        │  run() 中被引擎放入独立线程执行;next(key=..., text=...) 投喂数据
        ▼
pw.io.python.read(subject, schema=..., autocommit_duration_ms=1000)  → 实时 Table
        ▼
pw.io.csv.write(table, "output.csv")  每次提交将表的变更流写入 CSV
        ▼
pw.run() 阻塞运行,CTRL+C 触发 KeyboardInterrupt 退出

管线的四个组成部分都集中在同一个文件中,下面逐段拆解。

源码逐段解析

1. 引擎与许可证初始化

import pathway as pw

# To use advanced features with Pathway Live Data Framework Scale, get your free license key...
# To use Pathway Live Data Framework Community, comment out the line below.
pw.set_license_key("demo-license-key-with-telemetry")

BEARER_TOKEN = os.environ["TWITTER_API_TOKEN"]

脚本通过 pw.set_license_key(...) 声明运行模式:使用 Scale 高级功能时填入从官方获取的许可证;仅使用 Community(社区版)时按注释说明注释掉这一行即可。演示脚本中的 demo-license-key-with-telemetry 是一个带遥测的演示占位 key。Bearer Token 从环境变量 TWITTER_API_TOKEN 读取,而不是硬编码在文件里——这是示例刻意示范的安全做法。

2. TwitterClient:对接第三方 SDK,把回调翻译成 next

class TwitterClient(tweepy.StreamingClient):
    _subject: TwitterSubject

    def __init__(self, subject: TwitterSubject) -> None:
        super().__init__(BEARER_TOKEN)
        self._subject = subject

    def on_response(self, response) -> None:
        self._subject.next(
            key=response.data.id,
            text=response.data.text,
        )

TwitterClient 继承 tweepyStreamingClient,是纯粹的第三方 SDK 封装层,Pathway 对它的存在一无所知。关键在 on_response 回调:每当流式 API 推来一条推文,它就把它转换成符合后续 Schema 的关键字参数(keytext),再调用 self._subject.next(...) 投喂给 Subject。建议将"第三方回调 → 标准化的 Schema 列值"的翻译逻辑放在这一层,让 Subject 只关心 Pathway 侧的数据格式。

注意:该示例通过 next 传的是结构化字段key/text 两个命名列),而非把整条消息打包进单一的 data 列。这一点与 ConnectorSubject.next 的文档化示例 一致:next 按关键字参数发送一行,参数必须与传给 read 的 Schema 兼容。

3. TwitterSubject:连接器的核心实现

class TwitterSubject(pw.io.python.ConnectorSubject):
    _twitter_client: TwitterClient

    def __init__(self) -> None:
        super().__init__()
        self._twitter_client = TwitterClient(self)

    def run(self) -> None:
        self._twitter_client.sample()

    def on_stop(self) -> None:
        self._twitter_client.disconnect()

这是整篇文章的技术核心。ConnectorSubjectpython/pathway/io/python/init.py 中定义的一个抽象基类,只要求实现 run,其余全部是内嵌能力。示例用三个方法写清了连接器的完整契约:

方法 在示例中的职责 引擎何时调用
__init__ 实例化 TwitterClient 并把自己传给它,形成回调回路 构造阶段
run 调用 self._twitter_client.sample() 启动(近乎无限期的)推文流 由引擎在独立线程中启动(见源码 start()
on_stop 调用 disconnect() 优雅断开流式连接,做清理 run 结束或异常后(源码:finally 分支中先 on_stop()close()

start() 的实现 可以看到精确的生命周期语义:start 会启动一个守护线程执行 run,一旦 run 返回(有限连接器)或抛出异常,finally 里会先调用 on_stop() 做清理,再调用 close() 向缓冲区放入结束哨兵(FINISH_LITERAL)。所以对无限流而言,正确的中止路径是"外部中断 runon_stop 清理",示例通过 CTRL+C 触发 KeyboardInterrupt 走到这条路径。

4. Schema:把投喂数据声明为一张表

class InputSchema(pw.Schema):
    key: int = pw.column_definition(primary_key=True)
    text: str

Schema 把 Subject 通过 next 发出的关键字参数一一映射成表的列:key 被声明为主键列(primary_key=True),text 为普通文本列。主键的存在意味着引擎对同 key 行的处理遵循"按主键更新/去重"语义。关于 Schema 的完整定义方式(列类型、默认值、主键、数据类型),可参考仓库中配套教程 docs/2.developers/4.user-guide/20.connect/99.connectors/30.custom-python-connectors.md

5. 读入、落盘与运行

input = pw.io.python.read(
    TwitterSubject(),
    schema=InputSchema,
    autocommit_duration_ms=1000,
)

pw.io.csv.write(input, "output.csv")

try:
    pw.run()
except KeyboardInterrupt:
    print("Done.")

三个调用分别对应:接入(read)、输出(write)、运行(run)

  • pw.io.python.read(subject, schema=..., autocommit_duration_ms=1000):把 Subject 注册为数据源。其中 autocommit_duration_ms=1000 表示每 1000 毫秒自动提交一次累积的更新并推入数据流图;该参数在 read 的签名 中默认值为 1500 毫秒。提交语义决定了实时性的粒度——commit 之间的数据会被缓冲,只有提交后引擎才开始处理。
  • pw.io.csv.write(input, "output.csv"):把表的变更流(而非最终快照)持续写入 output.csv。如 python/pathway/io/csv/init.py 的文档示例所示,CSV 会在数据列之外追加 time(批次序号)与 diff1 表示新增、-1 表示删除)两列来描述增量变更。
  • try/except KeyboardInterruptpw.run() 阻塞运行整个数据流;用户按下 CTRL+C 后抛出 KeyboardInterrupt 并打印 Done.,这正是 README 第 5 步"按 CTRL+C 停止脚本"的代码体现。

ConnectorSubject 的运行机制:引擎与缓冲区

为了让上面的写法不只是"能跑",还需要理解它背后的机制。从源码 python/pathway/io/python/init.py 可以看到,每个 Subject 内部持有一个线程安全的 Queue 作为缓冲区,所有投喂方法最终都落到 self._buffer.put(...)

  • next(**kwargs):发送一行结构化数据,字段与 Schema 列一一对应(实现);
  • next_str(...) / next_bytes(...) / next_json(...):分别以 raw / binary / json 三种历史格式发送,通常落到单一的 data 列,在新代码中推荐直接用 next
  • commit():发送提交信号;配合 autocommit_duration_ms,决定"何时把缓冲的更新真正交还给引擎处理";
  • close():发送结束哨兵,声明不再有消息;run 正常结束时会被自动调用。

引擎侧通过 start/end/read/seek 等回调(封装进 api.PythonSubject)消费这个缓冲区。由此可以理解几个对本示例有用的推论:

  1. run 在独立线程中执行:阻塞式的 sample() 不会卡死引擎主循环,on_response 里的 next 只是写入缓冲区。
  2. 无限流不会自动 close:只有 run 返回(有限数据源)或出错时,finally 才触发 on_stop()close();无限流的中止完全依赖外部的 KeyboardInterrupt / disconnect
  3. 一个 Subject 对象只能被一个连接器使用一次:源码在 read 中做了守卫,重复使用会抛出 ValueError,需要创建两个独立对象。

运行步骤与环境准备

按照示例 README 的 5 步启动:

第 1 步:安装 Pathway

pip install pathway

安装方式细节(含可选扩展、Python 版本要求)见仓库内的安装教程 docs/2.developers/4.user-guide/10.introduction/20.installation.md。示例 README 中还提示:若要使用 Scale 高级特性,需要申请免费许可证 key 并填入脚本中的 pw.set_license_key(...);仅用 Community 则注释掉该行。

第 2 步:安装附加依赖

pip install -r requirements.txt

该示例所需的第三方库极简——requirements.txt 中只有一个固定版本依赖:

tweepy==4.13.0

第 3 步:提供 Twitter Bearer Token

export TWITTER_API_TOKEN=<BEARER_TOKEN>

Token 需要在 Twitter(X)开发者平台创建应用后获取,脚本通过 os.environ["TWITTER_API_TOKEN"] 读取。由于 Twitter 已关闭免费版流式 API,此步骤在实际执行时可能需要付费档的访问权限——这也再次说明本示例的价值在于模式本身,而非 Twitter 这个具体数据源

第 4 步:运行示例

python twitter_connector_example.py

脚本启动后,推文流被逐步写入 output.csv(含 time/diff 增量列),可以随时用另一个终端查看:

tail -f output.csv

第 5 步:停止脚本

# 在前台终端按下 CTRL+C

KeyboardInterrupt 被捕获后打印 Done. 并优雅退出。

代码仓库层面的行为验证

Python 连接器的行为在仓库测试中有充分覆盖,可作为理解本示例的佐证:

配套的官方教程 自定义 Python 连接器(custom Python connectors) 与本文档一一对应:先讲"读取静态文件"的最简有限场景,再讲本例"Twitter 无限流"的进阶场景,并系统列出 ConnectorSubject 需要实现的方法(runon_stop)与内嵌方法(nextnext_strnext_bytescommitclose)参考。

把 Twitter 换成你自己的数据源

README 强调:任何数据流都可以用这种方法喂给 Pathway。把示例抽象出来,一个自定义 Python 连接器只需满足一个最小骨架:

import pathway as pw

class MySubject(pw.io.python.ConnectorSubject):
    def run(self) -> None:
        # 从任意源(WebSocket、消息队列、轮询 HTTP、读文件、生成器等)取数,
        # 每拿到一条就调用一次 next,字段与下方 Schema 的列一一对应
        self.next(key=..., text=...)

class InputSchema(pw.Schema):
    key: int = pw.column_definition(primary_key=True)
    text: str

table = pw.io.python.read(
    MySubject(),
    schema=InputSchema,
    autocommit_duration_ms=1000,   # 每 1000ms 提交一批,可调小以获得更低延迟
)

pw.io.csv.write(table, "output.csv")
pw.run()

替换数据源时只需要修改 run 内的取数逻辑(以及 on_stop 里的清理逻辑),其余——Schema 声明、提交节奏、输出方式——完全复用。仓库中另一篇教程 从 Web 抓取新闻的自定义连接器 就是以这个思路对接网页抓取场景的又一实例。

使用建议与注意事项

  • 无限流的中止:不要在无限流里期待 run 自行返回,务必像示例一样用 on_stopdisconnect() 之类的清理,并依靠 CTRL+C/外部信号结束进程。
  • 提交节奏autocommit_duration_ms 是延迟与批量的折中——追求低延迟可调小,追求吞吐可调大;不传则采用源码中的默认值 1500 毫秒(见 read 签名)。
  • 大批量突发数据:若数据源会一次性爆发海量消息,read 还支持 max_backlog_size 参数限制积压、防止内存尖峰;设置后 Subject 内部的缓冲队列有界,投喂方法会在队列满时阻塞(详见 ConnectorSubject 类注释与 read 参数说明)。
  • Twitter API 现状:由于免费档流式 API 已关闭,运行本示例需要具备相应访问权限的 Token;把关注点放在"Subject + Schema + read + write"这套自定义接入模式上,才是本示例最值得沉淀的能力。
登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 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
531
594
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.36 K
1.46 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.01 K
516
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
547
388