无需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输出端口拉取数据。
步骤:
- 先安装依赖库:
pip install nifi-site-to-site - 在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用结构化流直接读取这些文件。
步骤:
- 在NiFi中配置
PutHDFS/PutS3Object处理器,将数据按规则输出为滚动文件(比如按小时分片,格式选Parquet)。 - PySpark端读取流式文件:
这个方案的优势是无需额外开发,依赖都是Spark和NiFi的原生功能,稳定性高。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()
方案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
相关产品推荐
相关产品推荐

