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

如何在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控制台生产者(快速测试)

  1. 在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
  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脚本(自动化/爬虫结果实时发送)

  1. 安装依赖:
pip install kafka-python
  1. 编写脚本(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实现,简单易维护:

  1. 在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数据后停止,适合每日调度场景。

  1. 提交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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 09:14:57