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

Structured Streaming+Kafka集成异常:numInputRows为0但offset上涨

问题描述

使用Spark 3.3版本进行Structured Streaming与Kafka集成时出现异常:初始阶段能正常获取正确的numInputRows,但约5次拉取后numInputRows变为0,而endOffset持续上涨。已尝试以下操作但问题仍未解决:

  • 删除checkpointLocation
  • 创建新Kafka Topic
  • 切换Spark版本至3.2.5

附上相关代码:

String s=" seq    , actual_voltage    , date    , set_current    , test_time    , set_voltage    , actual_current    , battery_temp1    , battery_temp2    , battery_temp3    , record_no    , actual_power    , setpower    , capacity    , energy    , relative_time    , step_num    , cycle_num    , status    , continuous_time    ,  step_discharge_cap    , step_discharge_energy    , accumulate_charge_cap    , accumulate_discharge_cap    , accumulate_charge_energy    , accumulate_discharge_energy    , max_auxiliary_voltage    , min_auxiliary_voltage    , max_auxiliary_temp    , min_auxiliary_temp    , auxiliary_voltage_diff    , auxiliary_temperature_diff    , total_auxiliary_voltage    , average_auxiliary_voltage    , average_auxiliary_temp    , max_cell_voltage    , min_cell_voltage    , max_cell_temp    , min_cell_temp    , cell_voltage_diff    , cell_temperature_diff    , average_cell_voltage    , average_cell_temp    , energy_judgment_value    , energy_judgment_std_value    , energy_judgment_q_value    , cap_judgment_value    , cap_judgment_std_value    , battery_max_temp    , battery_min_temp    , battery_avg_temp    , max_auxiliary_pressure    , min_auxiliary_pressure   ,pcNo  ,eventId , sendTime , topicCode ,pk ";
String[] sArray=s.split(",");
List<String> schemaList=new ArrayList<>();
for(int i=0;i<sArray.length;i++){
    schemaList.add(sArray[i].trim());
}

String schemaString=" seq  STRING , actual_voltage  STRING , date  STRING , set_current  STRING , test_time  STRING , set_voltage  STRING , actual_current  STRING , battery_temp1  STRING , battery_temp2  STRING , battery_temp3  STRING , record_no  STRING , actual_power  STRING , setpower  STRING , capacity  STRING , energy  STRING , relative_time  STRING , step_num  STRING , cycle_num  STRING , status  STRING , continuous_time  STRING ,  step_discharge_cap  STRING , step_discharge_energy  STRING , accumulate_charge_cap  STRING , accumulate_discharge_cap  STRING , accumulate_charge_energy  STRING , accumulate_discharge_energy  STRING , max_auxiliary_voltage  STRING , min_auxiliary_voltage  STRING , max_auxiliary_temp  STRING , min_auxiliary_temp  STRING , auxiliary_voltage_diff  STRING , auxiliary_temperature_diff  STRING , total_auxiliary_voltage  STRING , average_auxiliary_voltage  STRING , average_auxiliary_temp  STRING , max_cell_voltage  STRING , min_cell_voltage  STRING , max_cell_temp  STRING , min_cell_temp  STRING , cell_voltage_diff  STRING , cell_temperature_diff  STRING , average_cell_voltage  STRING , average_cell_temp  STRING , energy_judgment_value  STRING , energy_judgment_std_value  STRING , energy_judgment_q_value  STRING , cap_judgment_value  STRING , cap_judgment_std_value  STRING , battery_max_temp  STRING , battery_min_temp  STRING , battery_avg_temp  STRING , max_auxiliary_pressure  STRING , min_auxiliary_pressure  STRING ,pcNo STRING ,eventId STRING, sendTime STRING, topicCode STRING ,pk STRING";
ExpressionEncoder<Row> encoder = RowEncoder.apply(StructType.fromDDL(schemaString));

SparkSession sparkSession=SparkSession.builder()/*.master("local[*]")*/
        .config("spark.shuffle.useOldFetchProtocol","true")
        .config("spark.network.timeout","600s")
        .config("kafka.bootstrap.servers", "hadoop105:9092,hadoop104:9092")
        .config("spark.serializer","org.apache.spark.serializer.KryoSerializer").getOrCreate();
Dataset<Row> dataFrameReader =sparkSession.readStream().format("kafka")
        /*.option("startingOffsets", "earliest")*/
        /*.option("startingOffsets", "latest")*/
        .option("kafka.bootstrap.servers", "hadoop105:9092,hadoop104:9092")
        .option("kafka.fetch.max.wait.ms","600000")
        .option("kafka.max.partition.fetch.bytes","82428000")
        .option("kafka.fetch.message.max.bytes","82428800")
        .option("kafka.fetch.max.bytes","80000000")
        .option("kafka.receive.buffer.bytes","3000000")
        .option("maxOffsetsPerTrigger","500")
        .option("kafka.group.id","liu_test_group")
        .option("kafka.max.poll.records","500")
        .option("kafka.session.timeout.ms","120000")
        .option("subscribe", "hudi_test10_jlv9").load();

Dataset<Row> c11=dataFrameReader.selectExpr("CAST(value AS STRING)").as("a");

Dataset<TestingDataModal> c22=c11.map(new MapFunction<Row, TestingDataModal>() {
    @Override
    public TestingDataModal call(Row value) throws Exception {
        JSONObject json=JSON.parseObject(value.json());
        String jsonValue=json.get("value").toString();
        TestingDataModal testingDataModal=JSON.parseObject(jsonValue,TestingDataModal.class);
        testingDataModal.setId(IdUtil.fastUUID()+"_"+IdUtil.fastUUID());
        testingDataModal.setSendTime(new Date().getTime());
        return testingDataModal;
    }
},Encoders.bean(TestingDataModal.class));

Dataset<String> dataString=c22.flatMap(new FlatMapFunction<TestingDataModal, String>() {
    @Override
    public Iterator<String> call(TestingDataModal testingDataModal) throws Exception {
        String[][] tableValue=testingDataModal.getTableValue();
        String[] titleValue=testingDataModal.getTableTitle();
        List<String> valuesJson=new ArrayList<>();
        for(int i=0;i<tableValue.length;i++){
            JSONObject jsonRow=new JSONObject(true);
            for(int j=0;j<titleValue.length;j++){
                String key=titleValue[j];
                String value=tableValue[i][j];
                jsonRow.put(key,value);
            }

            jsonRow.put("pcNo",testingDataModal.getPcNo());
            jsonRow.put("eventId",testingDataModal.getEventId());
            jsonRow.put("sendTime",IdUtil.fastUUID());
            jsonRow.put("topicCode",testingDataModal.getTopicCode());
            jsonRow.put("pk", testingDataModal.getEventId()+"_"+tableValue[i][0]);
            valuesJson.add(jsonRow.toJSONString());
        }
        return valuesJson.iterator();
    }
},Encoders.STRING());

Dataset<Row> rowDataset=dataString.map(new MapFunction<String, Row>() {
    @Override
    public Row call(String value) throws Exception {
        JSONObject jsonObject1=JSONObject.parseObject(value);
        List<String> rowValue=new ArrayList<>();
        for(String schema:schemaList){
            Object o=jsonObject1.get(schema);
            if(o==null){
                rowValue.add("");
            }else {
                rowValue.add(o.toString());
            }
        }
        return RowFactory.create(rowValue.toArray());
    }
},encoder);

StreamingQuery query=rowDataset.writeStream().format("org.apache.hudi")
        .options(QuickstartUtils.getQuickstartWriteConfigs())
        .option(HoodieWriteConfig.PRECOMBINE_FIELD_NAME.key(), "sendTime")
        .option(KeyGeneratorOptions.RECORDKEY_FIELD_NAME.key(), "pk")
        .option(KeyGeneratorOptions.PARTITIONPATH_FIELD_NAME.key(), "eventId")
        .option(HoodieWriteConfig.TABLE_NAME, "datas")
        .option("hoodie.datasource.write.keygenerator.class","org.apache.hudi.keygen.ComplexKeyGenerator")
        .option(DataSourceWriteOptions.TABLE_TYPE().key(), HoodieTableType.MERGE_ON_READ.name())
        .option("write.retry.times","10")
        .option("checkpointLocation", "hdfs://hadoop101:8020/hudi_check5")
        .outputMode("append")
        .option("path", "hdfs://hadoop101:8020/hudi_5")
        .trigger(Trigger.ProcessingTime(30000))
        .start();

query.awaitTermination();
排查建议
  1. 修复Kafka消息解析逻辑
    原代码中c11.map阶段用value.json()解析Row会生成嵌套JSON结构,若原始消息格式不匹配会导致解析出空数据。替换为直接获取转换后的字符串:

    Dataset<TestingDataModal> c22 = c11.map(new MapFunction<Row, TestingDataModal>() {
        @Override
        public TestingDataModal call(Row value) throws Exception {
            String jsonValue = value.getString(0); // 直接提取STRING类型的value
            TestingDataModal testingDataModal = JSON.parseObject(jsonValue, TestingDataModal.class);
            testingDataModal.setId(IdUtil.fastUUID() + "_" + IdUtil.fastUUID());
            testingDataModal.setSendTime(new Date().getTime());
            return testingDataModal;
        }
    }, Encoders.bean(TestingDataModal.class));
    
  2. 定位数据丢失环节
    在每个转换阶段后添加临时控制台输出,确认数据在哪一步丢失:

    • 在c22后添加:c22.writeStream().format("console").outputMode("append").start();
    • 在dataString后添加:dataString.writeStream().format("console").outputMode("append").start();
      同时检查TestingDataModal的tableValue是否为空,若为空会导致flatMap无输出。
  3. 优化Kafka消费配置

    • 移除重复的kafka.bootstrap.servers配置(SparkSession和readStream中保留一个即可)
    • 取消maxOffsetsPerTrigger和kafka.max.poll.records的重复限制,避免触发间隔内数据拉取被过度约束
    • 开启Kafka消费DEBUG日志,查看是否有消息跳过或消费异常的日志记录
  4. 排查HUDI写入影响
    临时将输出改为控制台,确认rowDataset是否有数据输出:

    StreamingQuery query = rowDataset.writeStream()
            .format("console")
            .option("checkpointLocation", "hdfs://hadoop101:8020/hudi_check_temp")
            .outputMode("append")
            .trigger(Trigger.ProcessingTime(30000))
            .start();
    

    如果控制台有数据,说明问题出在HUDI环节:当前代码将sendTime设置为UUID字符串,而HUDI的PRECOMBINE_FIELD_NAME需要可比较的时间值,修正为时间戳:

    jsonRow.put("sendTime", new Date().getTime());
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 11:05:55