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

如何在Hortonworks中将Spark Streaming数据存储到HDFS?

将Spark Streaming数据从Kafka写入HDFS(Hortonworks环境)

我来帮你把现有代码改成将Kafka流数据写入HDFS,同时适配Hortonworks环境,一步步来:

一、环境前置检查

在动手改代码前,先确认Hortonworks集群的几个关键点:

  • 你的Spark集群已经配置好HDFS访问权限:Hortonworks默认Spark和HDFS是集成好的,只要运行Spark作业的用户拥有HDFS目标目录的写入权限就行。比如你要写到hdfs://nn01.example.com:8020/user/jayz/kafka-stream-output,先手动创建这个目录并赋权:
    hdfs dfs -mkdir -p /user/jayz/kafka-stream-output
    hdfs dfs -chmod 775 /user/jayz/kafka-stream-output
    
  • 确认Kafka和Spark版本兼容:Hortonworks的组件都是预适配的,只要你用的是集群自带的Spark依赖,就不会有版本问题。

二、修改Spark Streaming代码

核心是用DStream的saveAsTextFiles方法替代(或补充)控制台打印,以下是完整修改后的代码:

import _root_.kafka.serializer.DefaultDecoder
import _root_.kafka.serializer.StringDecoder
import org.apache.spark.SparkConf
import org.apache.spark.streaming.{Seconds, StreamingContext}
import org.apache.spark.streaming.kafka.KafkaUtils
import org.apache.spark.storage.StorageLevel

object StreamingDataNew {
  def main(args: Array[String]): Unit = {
    val sparkConf = new SparkConf().setAppName("Kafka-To-HDFS").setMaster("local[*]")
    // 如果你是在Hortonworks集群上运行,建议去掉setMaster("local[*]"),由spark-submit指定
    val ssc = new StreamingContext(sparkConf, Seconds(10))

    val kafkaConf = Map(
      "metadata.broker.list" -> "localhost:9092",
      "zookeeper.connect" -> "localhost:2181",
      "group.id" -> "kafka-streaming-example",
      "zookeeper.connection.timeout.ms" -> "200000"
    )

    val lines = KafkaUtils.createStream[Array[Byte], String, DefaultDecoder, StringDecoder](
      ssc,
      kafkaConf,
      Map("topic-one" -> 1), // 订阅topic-one,指定1个分区
      StorageLevel.MEMORY_ONLY
    )

    // 提取Kafka消息的value部分(你的原代码已经在做这个)
    val words = lines.flatMap { case (x, y) => y.split(" ") }

    // 1. 保留控制台打印用于调试
    words.print()

    // 2. 新增:将数据写入HDFS,替换成你的HDFS路径
    // 格式:前缀 + 自动生成的批次时间戳子目录
    val hdfsOutputPath = "hdfs://nn01.example.com:8020/user/jayz/kafka-stream-output"
    words.saveAsTextFiles(hdfsOutputPath)

    ssc.start()
    ssc.awaitTermination()
  }
}

代码修改关键点:

  1. HDFS路径配置:把hdfsOutputPath改成你集群实际的HDFS路径,Hortonworks的NameNode默认端口是8020,如果你不确定,可以在Ambari里查看HDFS的配置项dfs.namenode.http-address。
  2. saveAsTextFiles特性:这个方法会为每个微批次生成一个以时间戳命名的子目录,比如kafka-stream-output-1699999999000,里面是该批次的输出文件(默认是part-00000这类文件)。
  3. 集群运行注意:如果要在Hortonworks集群的YARN上运行,记得去掉代码里的setMaster("local[*]"),用spark-submit指定运行模式。

三、打包并运行作业

  1. 打包代码:用sbt或maven打包成jar包,确保依赖是provided(因为Hortonworks集群已经有Spark和Kafka的依赖)。
  2. 提交到Hortonworks集群:用spark-submit命令提交,比如:
    spark-submit \
      --class StreamingDataNew \
      --master yarn \
      --deploy-mode cluster \
      --executor-memory 2G \
      --num-executors 2 \
      your-spark-kafka-jar.jar
    
  3. 验证输出:运行后,你可以用HDFS命令查看输出:
    hdfs dfs -ls /user/jayz/kafka-stream-output-*
    hdfs dfs -cat /user/jayz/kafka-stream-output-1699999999000/part-00000
    

四、进阶优化(可选)

如果需要更灵活的输出(比如合并文件、自定义格式),可以用foreachRDD结合Hadoop的API,但saveAsTextFiles已经能满足基础的文本存储需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 09:22:40