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

发布时间:2026/9/8 20:23:26

Pathway 自定义 Python Connector 实战:用 `ConnectorSubject` 把 Twitter 等任意流式数据接入实时数据管道
Pathway 自定义 Python Connector 实战用ConnectorSubject把 Twitter 等任意流式数据接入实时数据管道【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway本文基于 Pathway 官方示例仓库中的 Twitter 自定义连接器examples/projects/custom-python-connector-twitter展开讲解如何通过继承pw.io.python.ConnectorSubject编写自己的 Python 输入连接器把外部数据源文中以 Twitter 流为例变成一张实时变化的 Pathway 表并持续写出到 CSV。读完本文你将掌握自定义 Python 连接器的生命周期回调、数据投喂方法、Schema 声明与批量提交语义并能把同一套模式复用到任意可被 Python 拉取的流式数据源上。示例背景与适用场景示例 READMEexamples/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_ms1000) → 实时 Table ▼ pw.io.csv.write(table, output.csv) 每次提交将表的变更流写入 CSV ▼ pw.run() 阻塞运行CTRLC 触发 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把回调翻译成nextclass 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( keyresponse.data.id, textresponse.data.text, )TwitterClient继承 tweepy 的StreamingClient是纯粹的第三方 SDK 封装层Pathway 对它的存在一无所知。关键在on_response回调每当流式 API 推来一条推文它就把它转换成符合后续 Schema 的关键字参数key、text再调用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()这是整篇文章的技术核心。ConnectorSubject是 python/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。所以对无限流而言正确的中止路径是外部中断run→on_stop清理示例通过CTRLC触发 KeyboardInterrupt 走到这条路径。4. Schema把投喂数据声明为一张表class InputSchema(pw.Schema): key: int pw.column_definition(primary_keyTrue) text: strSchema 把 Subject 通过next发出的关键字参数一一映射成表的列key被声明为主键列primary_keyTruetext为普通文本列。主键的存在意味着引擎对同 key 行的处理遵循按主键更新/去重语义。关于 Schema 的完整定义方式列类型、默认值、主键、数据类型可参考仓库中配套教程 docs/2.developers/4.user-guide/20.connect/99.connectors/30.custom-python-connectors.md。5. 读入、落盘与运行input pw.io.python.read( TwitterSubject(), schemaInputSchema, autocommit_duration_ms1000, ) 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_ms1000)把 Subject 注册为数据源。其中autocommit_duration_ms1000表示每 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()阻塞运行整个数据流用户按下CTRLC后抛出KeyboardInterrupt并打印Done.这正是 README 第 5 步按 CTRLC 停止脚本的代码体现。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列在新代码中推荐直接用nextcommit()发送提交信号配合autocommit_duration_ms决定何时把缓冲的更新真正交还给引擎处理close()发送结束哨兵声明不再有消息run正常结束时会被自动调用。引擎侧通过 start/end/read/seek 等回调封装进api.PythonSubject消费这个缓冲区。由此可以理解几个对本示例有用的推论run在独立线程中执行阻塞式的sample()不会卡死引擎主循环on_response里的next只是写入缓冲区。无限流不会自动close只有run返回有限数据源或出错时finally才触发on_stop()与close()无限流的中止完全依赖外部的KeyboardInterrupt/disconnect。一个 Subject 对象只能被一个连接器使用一次源码在 read 中做了守卫重复使用会抛出ValueError需要创建两个独立对象。运行步骤与环境准备按照示例 README 的 5 步启动第 1 步安装 Pathwaypip 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 中只有一个固定版本依赖tweepy4.13.0第 3 步提供 Twitter Bearer Tokenexport TWITTER_API_TOKENBEARER_TOKENToken 需要在 TwitterX开发者平台创建应用后获取脚本通过os.environ[TWITTER_API_TOKEN]读取。由于 Twitter 已关闭免费版流式 API此步骤在实际执行时可能需要付费档的访问权限——这也再次说明本示例的价值在于模式本身而非 Twitter 这个具体数据源。第 4 步运行示例python twitter_connector_example.py脚本启动后推文流被逐步写入output.csv含time/diff增量列可以随时用另一个终端查看tail -f output.csv第 5 步停止脚本# 在前台终端按下 CTRLCKeyboardInterrupt被捕获后打印Done.并优雅退出。代码仓库层面的行为验证Python 连接器的行为在仓库测试中有充分覆盖可作为理解本示例的佐证python/pathway/tests/test_io.py 中的 test_python_connector演示了next_json/next_str/next_bytes三种投喂方式产生的行最终在表中等价于手写数据验证了回调里发多少条、表里就有多少条的语义同一文件中的 test_python_connector_on_stop在run为空实现立即返回的有限连接器上断言on_stop一定被引擎调用——印证上文介绍的run结束后必然清理生命周期。配套的官方教程 自定义 Python 连接器custom Python connectors 与本文档一一对应先讲读取静态文件的最简有限场景再讲本例Twitter 无限流的进阶场景并系统列出ConnectorSubject需要实现的方法run、on_stop与内嵌方法next、next_str、next_bytes、commit、close参考。把 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_keyTrue) text: str table pw.io.python.read( MySubject(), schemaInputSchema, autocommit_duration_ms1000, # 每 1000ms 提交一批可调小以获得更低延迟 ) pw.io.csv.write(table, output.csv) pw.run()替换数据源时只需要修改run内的取数逻辑以及on_stop里的清理逻辑其余——Schema 声明、提交节奏、输出方式——完全复用。仓库中另一篇教程 从 Web 抓取新闻的自定义连接器 就是以这个思路对接网页抓取场景的又一实例。使用建议与注意事项无限流的中止不要在无限流里期待run自行返回务必像示例一样用on_stop做disconnect()之类的清理并依靠CTRLC/外部信号结束进程。提交节奏autocommit_duration_ms是延迟与批量的折中——追求低延迟可调小追求吞吐可调大不传则采用源码中的默认值1500毫秒见 read 签名。大批量突发数据若数据源会一次性爆发海量消息read还支持max_backlog_size参数限制积压、防止内存尖峰设置后 Subject 内部的缓冲队列有界投喂方法会在队列满时阻塞详见 ConnectorSubject 类注释与 read 参数说明。Twitter API 现状由于免费档流式 API 已关闭运行本示例需要具备相应访问权限的 Token把关注点放在Subject Schema read write这套自定义接入模式上才是本示例最值得沉淀的能力。【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关新闻

res-downloader 快速上手指南

res-downloader 快速上手指南

2026/9/8 20:23:26

res-downloader 快速上手指南 【免费下载链接】res-downloader 视频号、小程序、抖音、快手、小红书、直播流、m3u8、酷狗、QQ音乐等常见网络资源下载! 项目地址: https://gitcode.com/GitHub_Trending/re/res-downloader 你有没有遇到过:视频号里一个特别好…

LangChain核心三件事:链式抽象、Agent组件与LangGraph编排

LangChain核心三件事:链式抽象、Agent组件与LangGraph编排

2026/9/8 20:13:26

这两年做 AI 应用开发,如果你还没碰过 LangChain 生态,那确实有点说不过去了。从最简单的“把 Prompt 接到大模型上”到复杂的多步骤工作流、智能体系统,LangChain 几乎把开发过程中的方方面面都抽象成了可组合的模块。但问题也出在这&#x…

论文盲审意见怎么应对?修改策略一篇讲透

论文盲审意见怎么应对?修改策略一篇讲透

2026/9/8 20:13:26

拿到盲审意见的那一刻,很多人会先紧张:措辞犀利,结论栏又不明朗,完全不知道怎么应对。这篇把盲审返修拆成「读懂意见—定位修改点—逐条落实—复查回稿」四个阶段,每个阶段给出可执行的修改策略,帮你把「不…

ESP32-S3端云协同AI架构:从语音唤醒到自主演进的工程实践

ESP32-S3端云协同AI架构:从语音唤醒到自主演进的工程实践

2026/9/8 21:03:28

1. 项目概述:一块开发板如何长出“感知-思考-表达”的神经网络 你手边那块不到百元的 ESP32-S3 开发板,表面看只是个带双核 Xtensa LX7、2.4GHz Wi-Fi Bluetooth LE、USB OTG 和丰富外设接口的微控制器——但它的真正价值,从来不在参数表里&…

从陶瓷工业百强看京尚“市场与品质双轮驱动”的实战逻辑

从陶瓷工业百强看京尚“市场与品质双轮驱动”的实战逻辑

2026/9/8 21:03:28

前段时间陶瓷行业圈子里最热闹的一件事,就是新一届全国陶瓷工业百强名单出炉。京尚这个品牌不仅稳稳上榜,还成了榜单里被反复提及的“双轮驱动”典型——市场和品质两头都抓得硬。我做这行十几年,见过太多企业要么拼命冲销量把品质丢了&#…

tiktoken 分词器完整指南:如何为 OpenAI 模型精确计算 token

tiktoken 分词器完整指南:如何为 OpenAI 模型精确计算 token

2026/9/8 21:03:28

tiktoken 分词器完整指南:如何为 OpenAI 模型精确计算 token 【免费下载链接】tiktoken tiktoken is a fast BPE tokeniser for use with OpenAIs models. 项目地址: https://gitcode.com/GitHub_Trending/ti/tiktoken 调用 OpenAI API 前,你需要…

ROS2 Launch 文件完全指南:从手动多终端到一键启动与参数化复用

ROS2 Launch 文件完全指南:从手动多终端到一键启动与参数化复用

2026/9/8 21:03:28

1. 为什么你需要 Launch:从手动开终端的痛说起1.1 一个过来人脑中的"标准化痛苦"刚开始接触 ROS2 的人,大多经历过这样一段蹒跚期:装好了 Humble,跟着教程敲ros2 run turtlesim turtlesim_node,小乌龟出来了…

渔业目标检测数据集使用指南:标注诊断与YOLO实战

渔业目标检测数据集使用指南:标注诊断与YOLO实战

2026/9/8 21:03:28

简介:本资源是面向计算机视觉初学者与AI安全监控开发者的小型钓鱼行为检测专用数据集,聚焦于岸边钓鱼人员的识别与定位任务,适用于智能公园管理、水域保护及安防预警等实际场景。压缩包共2000个文件,含1000张JPG格式原始图像与100…

15 分钟跑通 RPCS3:PS3 模拟器的三个落地场景——跑游戏、打补丁、调崩溃

15 分钟跑通 RPCS3:PS3 模拟器的三个落地场景——跑游戏、打补丁、调崩溃

2026/9/8 20:53:28

15 分钟跑通 RPCS3:PS3 模拟器的三个落地场景——跑游戏、打补丁、调崩溃 【免费下载链接】rpcs3 PlayStation 3 emulator and debugger 项目地址: https://gitcode.com/GitHub_Trending/rp/rpcs3 想让 PS3 光盘游戏在电脑上跑起来,还能在崩溃时定…

中国人民大学杨琳团队《Nature Communications》 | 全球潮汐湿地土壤有机碳时空格局与环境驱动:一项2009-2020年的全球评估

中国人民大学杨琳团队《Nature Communications》 | 全球潮汐湿地土壤有机碳时空格局与环境驱动:一项2009-2020年的全球评估

2026/9/7 20:21:46

本文首发于“生态学者”!从“湿地面积”到“土壤碳密度”:为什么需要重新认识潮汐湿地蓝碳变化?潮汐湿地位于陆地与海洋的交汇地带,包括红树林、盐沼和潮滩,是全球重要的蓝碳生态系统。其土壤能够长期储存大量有机碳&a…

adb抓包

adb抓包

2026/9/8 4:55:53

前言 本文介绍如何通过 tcpdump 在 Android 手机上抓取网络数据包,并在电脑端使用 Wireshark 进行分析。适用于需要排查 App 网络请求、分析接口调用或调试网络问题的开发与测试场景。1. 手机要有 root 权限2. 下载 tcpdump3. adb push C:\Users\zhangkuixun\Downlo…

大模型推理镜像极简瘦身:从 25GB 巨无霸到 3GB 精简镜像实战

大模型推理镜像极简瘦身:从 25GB 巨无霸到 3GB 精简镜像实战

2026/9/7 8:03:37

大模型推理镜像极简瘦身:从 25GB 巨无霸到 3GB 精简镜像实战 在云原生基础设施中,容器镜像体积直接决定了服务的部署速度与弹性扩容敏捷度。对于传统的 Go / Java 微服务,镜像体积通常被严格控制在 50MB 到 200MB 以内,拉取镜像只…

芯片良率波动可视化:动画拆解工艺因果,重建客户信任

芯片良率波动可视化:动画拆解工艺因果,重建客户信任

2026/9/8 0:02:30

芯片这个行业有个不太被人摆到台面上、但几乎每天都在发生的场景:客户拿着一条良率曲线截图问你,这批货的良率怎么掉了三个点,是不是工艺出问题了,产生的不良会不会流到他们产线上去。你解释了半天,客户似懂非懂&#…

PyTorch DataLoader参数冲突:sampler与shuffle互斥的根源与正确写法

PyTorch DataLoader参数冲突:sampler与shuffle互斥的根源与正确写法

2026/9/8 0:02:30

ValueError: sampler option is mutually exclusive with shuffle,这个报错我在 PyTorch 的 DataLoader 上至少见过几十次了,而且很有意思的是,它经常不是新手专属——很多写了好几年模型的老手,在从单机改成自定义采样器&#xf…

中国车企再破谣言,GAC吉利零跑获欧盟安全五星

中国车企再破谣言,GAC吉利零跑获欧盟安全五星

2026/9/8 0:02:30

有人可能在网上开着皮卡拍视频,声称中国电动车不仅性能不如美国大排量车型,安全性也堪忧。然而事实恰恰相反,GAC、吉利和零跑最新推出的电动车型在极为严苛的欧盟新车安全评鉴(Euro NCAP)测试中全部斩获满分。就在特斯…

远程协作的工作台整理

远程协作的工作台整理

2026/9/8 4:23:39

远程协作的工作台整理远程协作的核心不是再加一个工具,而是让交接信息足够完整。异步任务要写明目标、输入位置、完成标准和需要决策的人。 工作台的最小配置 将日程、待办、代码和沟通入口收拢到少数固定位置;通知按紧急程度分层。工作台不需要模仿办公…

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

2026/9/8 3:19:39

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

2026/9/8 4:00:23

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…