PyFlink自定义SourceFunction求助:定时调用WebAPI示例
Python实现Flink自定义SourceFunction(每分钟调用WebAPI)
PyFlink完全支持自定义SourceFunction实现,你的示例代码存在构造方法调用错误,且缺少定时调用WebAPI的逻辑,下面给出完整可运行的实现方案。
问题分析
你的代码中super().__init__(self)是错误写法——Python调用父类构造方法无需传递self参数,应改为super().__init__()。同时代码未实现每分钟调用的定时逻辑,也没有WebAPI请求相关处理。
完整实现示例
以下代码实现了每分钟调用一次WebAPI的自定义Source,包含异常处理与任务取消逻辑:
from pyflink.datastream import SourceFunction import requests import time from typing import Optional class WebApiSource(SourceFunction): def __init__(self, api_url: str): super().__init__() self.api_url = api_url self._running = True def run(self, ctx: SourceFunction.SourceContext): while self._running: try: # 调用WebAPI并解析响应 response = requests.get(self.api_url) response.raise_for_status() api_data = response.json() # 将数据发送到Flink数据流 ctx.collect(api_data) # 等待60秒,实现每分钟调用一次 time.sleep(60) except Exception as e: # 捕获异常避免任务崩溃,可根据需求添加日志逻辑 print(f"API调用失败: {str(e)}") time.sleep(10) # 出错后等待10秒重试 def cancel(self): self._running = False
作业使用示例
在Flink流作业中集成该自定义Source:
from pyflink.datastream import StreamExecutionEnvironment def main(): env = StreamExecutionEnvironment.get_execution_environment() # 设置并行度为1,避免多实例同时调用API(可根据业务需求调整) env.set_parallelism(1) # 替换为实际的WebAPI地址 custom_source = WebApiSource(api_url="https://your-api-address.com/data") data_stream = env.add_source(custom_source) # 打印输出,可替换为业务处理逻辑 data_stream.print() env.execute("WebAPI Custom Source Job") if __name__ == "__main__": main()
注意事项
- 确保PyFlink环境已安装
requests依赖:执行pip install requests - 若API需要认证、请求参数,可在
run方法中补充对应逻辑 time.sleep(60)是基础定时方式,如需更精确调度,可结合datetime模块计算下次执行时间- 生产环境建议使用
logging模块记录日志,替代print输出 - 若需并行调用API,可调整并行度,但需确认API支持多并发请求
内容的提问来源于stack exchange,提问作者nagendrasri
相关产品推荐
相关产品推荐

