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();
排查建议
修复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));定位数据丢失环节
在每个转换阶段后添加临时控制台输出,确认数据在哪一步丢失:- 在
c22后添加:c22.writeStream().format("console").outputMode("append").start(); - 在
dataString后添加:dataString.writeStream().format("console").outputMode("append").start();
同时检查TestingDataModal的tableValue是否为空,若为空会导致flatMap无输出。
- 在
优化Kafka消费配置
- 移除重复的
kafka.bootstrap.servers配置(SparkSession和readStream中保留一个即可) - 取消
maxOffsetsPerTrigger和kafka.max.poll.records的重复限制,避免触发间隔内数据拉取被过度约束 - 开启Kafka消费DEBUG日志,查看是否有消息跳过或消费异常的日志记录
- 移除重复的
排查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
相关产品推荐
相关产品推荐

