如何在HDP沙箱中通过Kafka+Oozie实现CSV数据导入及爬虫自动存HDFS?
解决方案:Ubuntu本地数据(CSV/爬虫结果)经Kafka导入HDP HDFS + Oozie每日调度
一、Ubuntu本地CSV发送到Kafka Topic
先确保HDP沙箱的Kafka服务可被Ubuntu访问:打开沙箱防火墙9092端口(默认Kafka端口),或修改Kafka配置advertised.listeners为沙箱IP:9092后重启Kafka。
方法1:Kafka控制台生产者(快速测试)
- 在HDP沙箱创建目标Topic:
/usr/hdp/current/kafka-broker/bin/kafka-topics.sh --create --topic hdfs_data_topic --bootstrap-server <沙箱IP>:9092 --partitions 2 --replication-factor 1
- 在Ubuntu本地执行(需有对应HDP版本的Kafka客户端):
kafka-console-producer.sh --broker-list <沙箱IP>:9092 --topic hdfs_data_topic < /path/to/your/local/file.csv
注:该方式会将CSV每行作为一条Kafka消息,适配多数场景需求。
方法2:Python脚本(自动化/爬虫结果实时发送)
- 安装依赖:
pip install kafka-python
- 编写脚本(
kafka_producer.py),支持CSV读取或直接发送爬虫结果:
from kafka import KafkaProducer import csv # 初始化Kafka生产者 producer = KafkaProducer(bootstrap_servers=['<沙箱IP>:9092'], value_serializer=lambda x: x.encode('utf-8')) # 读取CSV并发送 with open('/path/to/local/file.csv', 'r') as f: reader = csv.reader(f) next(reader) # 按需跳过表头 for row in reader: message = ','.join(row) producer.send('hdfs_data_topic', value=message) producer.flush() # 爬虫结果直接发送示例:producer.send('hdfs_data_topic', value=scraped_content)
执行脚本:python3 kafka_producer.py
二、从Kafka消费数据写入HDP HDFS
用HDP自带的Spark Structured Streaming实现,简单易维护:
- 在HDP沙箱编写Spark脚本(
kafka_to_hdfs.py):
from pyspark.sql import SparkSession from pyspark.sql.functions import col spark = SparkSession.builder.appName("KafkaToHDFS").getOrCreate() # 订阅Kafka Topic df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "<沙箱IP>:9092") \ .option("subscribe", "hdfs_data_topic") \ .load() # 解析Kafka消息为字符串 parsed_df = df.select(col("value").cast("string").alias("content")) # 写入HDFS(按日期分区,适配每日调度) query = parsed_df.writeStream \ .outputMode("append") \ .format("csv") \ .option("path", "/user/hdfs/crawled_data") \ .option("checkpointLocation", "/user/hdfs/checkpoint/kafka_to_hdfs") \ .option("header", "true") \ .partitionBy("date") \ .trigger(once=True) # 单次执行,适配Oozie调度 .start() query.awaitTermination()
注:
trigger(once=True)表示处理完当前Topic数据后停止,适合每日调度场景。
- 提交Spark作业:
spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.11:<对应HDP Spark版本> kafka_to_hdfs.py
替换
<对应HDP Spark版本>为你的HDP Spark版本,比如2.4.7.3.1.5.0-152。
三、Oozie实现每日自动爬取+数据流转
1. 准备Oozie工作流文件
创建workflow.xml:
<workflow-app xmlns="uri:oozie:workflow:0.5" name="data-pipeline-workflow"> <start to="run-crawler-producer"/> <action name="run-crawler-producer"> <ssh xmlns="uri:oozie:ssh-action:0.2"> <host><Ubuntu主机IP></host> <command>python3 /home/ubuntu/crawler_and_producer.py</command> <user>ubuntu</user> <password><Ubuntu用户密码></password> </ssh> <ok to="run-spark-consumer"/> <error to="fail"/> </action> <action name="run-spark-consumer"> <spark xmlns="uri:oozie:spark-action:0.2"> <job-tracker>${jobTracker}</job-tracker> <name-node>${nameNode}</name-node> <master>yarn</master> <mode>cluster</mode> <jar>/path/to/kafka_to_hdfs.py</jar> <spark-opts>--packages org.apache.spark:spark-sql-kafka-0-10_2.11:2.4.7.3.1.5.0-152</spark-opts> </spark> <ok to="end"/> <error to="fail"/> </action> <kill name="fail"> <message>Workflow failed, error message[${wf:errorMessage(wf:lastErrorNode())}]</message> </kill> <end name="end"/> </workflow-app>
说明:
run-crawler-producer通过SSH调用Ubuntu本地合并后的爬虫+Kafka生产脚本;run-spark-consumer执行Spark作业将Kafka数据写入HDFS。
2. 配置Oozie协调器(定时调度)
创建coordinator.xml,设置每日凌晨2点执行:
<coordinator-app xmlns="uri:oozie:coordinator:0.4" name="daily-data-pipeline" frequency="${coord:days(1)}" start="2024-01-01T02:00Z" end="2025-12-31T02:00Z" timezone="Asia/Shanghai"> <action> <workflow> <app-path>${nameNode}/user/hdfs/oozie/workflows/data-pipeline</app-path> <configuration> <property> <name>jobTracker</name> <value>${jobTracker}</value> </property> <property> <name>nameNode</name> <value>${nameNode}</value> </property> </configuration> </workflow> </action> </coordinator-app>
3. 部署并启动Oozie任务
- 将工作流文件、脚本上传到HDFS对应目录,比如
/user/hdfs/oozie/workflows/data-pipeline/ - 启动Oozie协调器:
oozie job -oozie http://<沙箱IP>:11000/oozie -config job.properties -run
job.properties需配置nameNode、jobTracker等HDP集群参数。
内容的提问来源于stack exchange,提问作者Xedonedron
相关产品推荐
相关产品推荐

