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

Java环境下使用Spark Streaming从HDFS读取JSON文件的技术求助

我来帮你把这个Spark Streaming读取HDFS单行JSON文件的实现梳理清楚,下面是完整的可运行代码和关键注意点:

完整实现代码

首先导入必要的依赖包,然后编写主逻辑:

import org.apache.spark.SparkConf;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.streaming.Durations;
import org.apache.spark.streaming.api.java.JavaDStream;
import org.apache.spark.streaming.api.java.JavaStreamingContext;
import org.apache.hadoop.mapreduce.lib.input.TextInputFormat;

public class HdfsJsonStreamingJob {
    public static void main(String[] args) throws InterruptedException {
        // 初始化Spark配置和Streaming上下文
        SparkConf config = new SparkConf().setAppName("HDFS Streaming Job").setMaster("local[*]");
        JavaStreamingContext jssc = new JavaStreamingContext(config, Durations.seconds(5));

        // 配置要监控的HDFS目录
        String hdfsTargetDir = "hdfs://your-nn-host:9000/path/to/your/json-folder";
        
        // 读取HDFS目录下的新增文件,提取单行JSON字符串
        JavaDStream<String> jsonLineStream = jssc.fileStream(
                hdfsTargetDir,
                String.class,
                String.class,
                TextInputFormat.class
        ).map(fileContentTuple -> fileContentTuple._2());

        // 将DStream转换为Dataset进行JSON解析与处理
        jsonLineStream.foreachRDD(jsonRdd -> {
            // 避免空RDD导致的解析异常
            if (!jsonRdd.isEmpty()) {
                // 获取或复用SparkSession实例
                SparkSession spark = SparkSession.builder()
                        .config(jsonRdd.sparkContext().getConf())
                        .getOrCreate();
                
                // 解析单行JSON为Dataset<Row>
                Dataset<Row> jsonDataset = spark.read().json(jsonRdd);

                // 这里替换为你的业务逻辑,比如打印Schema、数据预览、写入存储等
                System.out.println("JSON数据Schema:");
                jsonDataset.printSchema();
                System.out.println("前5条数据内容:");
                jsonDataset.show(5, false);
            }
        });

        // 启动Streaming作业并等待终止
        jssc.start();
        jssc.awaitTermination();
    }
}
关键注意事项
  • 文件写入规范:往监控的HDFS目录写文件时,一定要用原子性写入(比如先写到临时目录,再rename到目标目录),否则Streaming可能读取到未写完的不完整文件。
  • SparkSession复用:在foreachRDD内部创建SparkSession时,必须通过当前RDD的Spark上下文配置来构建,避免重复创建Session导致资源浪费。
  • 空RDD判断:监控周期内没有新增文件时,RDD会是空的,必须先判断非空再执行解析,否则会抛出异常。
  • JSON格式匹配:你的文件是单行JSON格式,正好匹配spark.read().json()的默认行为(每行视为一个独立JSON记录),不需要额外配置。
  • 依赖版本匹配:确保Spark Streaming和Spark SQL的依赖版本一致,比如都使用3.3.0版本,避免兼容性问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:25:00