You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.07 23:26:11