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

无需Kafka:PySpark与NiFi直接集成实现流数据读取

PySpark直接读取NiFi输出端口流数据(无需Kafka中转)

我完全理解你的需求——想用PySpark直接从NiFi输出端口读取流数据,绕开Kafka中转。确实Scala那边有官方的SiteToSiteClient可以直接对接,但Python生态里没有官方等价模块,不过有几个实用的方案可以解决这个问题:

方案1:利用第三方Python库实现NiFi Site-to-Site通信

社区开发了适配NiFi Site-to-Site协议的Python库nifi-site-to-site,可以模拟Scala版本SiteToSiteClient的核心功能,直接从NiFi输出端口拉取数据。

步骤:

  1. 先安装依赖库:
    pip install nifi-site-to-site
    
  2. 在PySpark代码中初始化客户端并拉取数据:
    from nifi_site_to_site import NiFiClient
    from pyspark.sql import SparkSession
    
    spark = SparkSession.builder.appName("NiFiDirectRead").getOrCreate()
    sc = spark.sparkContext
    
    # 初始化NiFi客户端
    nifi_client = NiFiClient(
        nifi_host="你的NiFi主机地址",
        nifi_port=8080,  # NiFi默认HTTP端口,根据实际配置修改
        output_port_name="目标输出端口名称",
        protocol="HTTP"  # 可选RAW,需和NiFi Site-to-Site配置一致
    )
    
    # 持续拉取并处理数据(适配流式场景)
    while True:
        # 拉取批量流文件
        flow_files = nifi_client.list_flow_files()
        for flow_file in flow_files:
            # 读取流文件内容(假设是JSON格式,可根据实际数据调整)
            content_str = flow_file.content.decode("utf-8")
            # 转换为Spark DataFrame
            df = spark.read.json(sc.parallelize([content_str]))
            # 这里添加你的数据处理逻辑
            df.show()
        # 可添加适当的休眠时间,避免过度请求
        import time
        time.sleep(5)
    

    注意:需要确保NiFi已开启Site-to-Site服务,并且输出端口允许该客户端访问;如果是生产环境,建议处理数据消费确认,避免重复读取。

方案2:NiFi输出到文件系统,PySpark读取流式文件

如果不想引入第三方库,这个方案更简单可靠:让NiFi把数据输出到HDFS、S3或本地文件系统的滚动分片文件(比如按时间生成的Parquet/JSON文件),然后PySpark用结构化流直接读取这些文件。

步骤:

  1. 在NiFi中配置PutHDFS/PutS3Object处理器,将数据按规则输出为滚动文件(比如按小时分片,格式选Parquet)。
  2. PySpark端读取流式文件:
    from pyspark.sql import SparkSession
    
    spark = SparkSession.builder.appName("NiFiFileStreamRead").getOrCreate()
    
    # 读取NiFi输出的流式文件
    stream_df = spark.readStream \
        .format("parquet") \
        .option("path", "/nifi/output/hdfs/path") \
        .option("maxFilesPerTrigger", 2)  # 每次触发读取的文件数
        .load()
    
    # 数据处理与输出示例
    query = stream_df.writeStream \
        .format("console") \
        .outputMode("append") \
        .start()
    
    query.awaitTermination()
    
    这个方案的优势是无需额外开发,依赖都是Spark和NiFi的原生功能,稳定性高。

方案3:调用NiFi REST API拉取数据

NiFi提供了完整的REST API,可以直接调用接口获取输出端口的流文件数据。这种方式灵活性高,但需要自己处理数据的消费确认和流式逻辑。

核心示例代码:

import requests
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("NiFiRESTRead").getOrCreate()
sc = spark.sparkContext

# NiFi REST API基础配置
nifi_base_url = "http://你的NiFi主机:8080/nifi-api"
queue_id = "目标输出端口的队列ID"  # 可通过NiFi UI获取队列ID
headers = {"Content-Type": "application/json"}

def fetch_nifi_data():
    # 拉取队列中的流文件
    fetch_url = f"{nifi_base_url}/flowfile-queues/{queue_id}/flowfiles"
    response = requests.post(fetch_url, headers=headers, json={"count": 10})
    if response.status_code == 200:
        flow_files = response.json()["flowFiles"]
        data_list = []
        for ff in flow_files:
            # 获取流文件内容
            content_url = f"{nifi_base_url}/flowfile-queues/{queue_id}/flowfiles/{ff['uuid']}/content"
            content_response = requests.get(content_url)
            data_list.append(content_response.text)
            # 确认消费,删除已处理的流文件
            delete_url = f"{nifi_base_url}/flowfile-queues/{queue_id}/flowfiles/{ff['uuid']}"
            requests.delete(delete_url, headers=headers)
        return data_list
    return []

# 流式处理循环
while True:
    data = fetch_nifi_data()
    if data:
        df = spark.read.json(sc.parallelize(data))
        df.show()
    import time
    time.sleep(5)

注意:生产环境中需要处理API认证、异常重试等逻辑,避免数据丢失或重复处理。


内容的提问来源于stack exchange,提问作者Atanu chatterjee

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 12:19:37