如何在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() } }
代码修改关键点:
- HDFS路径配置:把
hdfsOutputPath改成你集群实际的HDFS路径,Hortonworks的NameNode默认端口是8020,如果你不确定,可以在Ambari里查看HDFS的配置项dfs.namenode.http-address。 - saveAsTextFiles特性:这个方法会为每个微批次生成一个以时间戳命名的子目录,比如
kafka-stream-output-1699999999000,里面是该批次的输出文件(默认是part-00000这类文件)。 - 集群运行注意:如果要在Hortonworks集群的YARN上运行,记得去掉代码里的
setMaster("local[*]"),用spark-submit指定运行模式。
三、打包并运行作业
- 打包代码:用sbt或maven打包成jar包,确保依赖是
provided(因为Hortonworks集群已经有Spark和Kafka的依赖)。 - 提交到Hortonworks集群:用
spark-submit命令提交,比如:spark-submit \ --class StreamingDataNew \ --master yarn \ --deploy-mode cluster \ --executor-memory 2G \ --num-executors 2 \ your-spark-kafka-jar.jar - 验证输出:运行后,你可以用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
相关产品推荐
相关产品推荐

