如何手动创建KafkaRecord<String, GenericRecord>用于Beam单元测试及排障
一、手动创建KafkaRecord<String, GenericRecord>用于单元测试
KafkaRecord是Apache Beam的公开类,可直接实例化用于单元测试,步骤如下:
构建GenericRecord实例
使用Avro的GenericData.Record或GenericRecordBuilder基于目标Schema创建:// 示例Schema,替换为你的实际Schema Schema valueSchema = new Schema.Parser().parse(""" { "type": "record", "name": "User", "fields": [{"name": "id", "type": "int"}, {"name": "name", "type": "string"}] } """); GenericRecord valueRecord = new GenericData.Record(valueSchema); valueRecord.put("id", 1001); valueRecord.put("name", "test-user");实例化KafkaRecord
传入topic、分区、偏移量、时间戳、key、value等必要参数:KafkaRecord<String, GenericRecord> kafkaRecord = new KafkaRecord<>( "your-topic-name", 0, // 分区号 1L, // 偏移量 System.currentTimeMillis(), // 时间戳 "test-key", // Kafka key valueRecord, // Kafka value Collections.emptyList() // Kafka Headers,可选 );在测试Pipeline中使用
用Create转换生成输入PCollection,并指定正确的Coder:PCollection<KafkaRecord<String, GenericRecord>> input = testPipeline .apply(Create.of(Collections.singletonList(kafkaRecord))) .setCoder(KafkaRecordCoder.of( StringUtf8Coder.of(), GenericRecordCoder.of(valueSchema) ));
二、PAssert验证KV<GenericRecord, GenericRecord>时的EOFException排查
错误原因分析
从栈轨迹看,Direct Runner在克隆数据时,GenericRecord的Coder解码失败触发EOFException。核心问题是Coder配置不完整或测试数据与Pipeline输出的Schema不匹配。
解决步骤
为GenericRecordCoder指定明确Schema
原代码中GenericRecordCoder.of()未传入Schema,Avro无法正确编码/解码。需替换为对应字段的实际Schema:// 替换原setCoder代码,传入key和value的Schema KvCoder<GenericRecord, GenericRecord> kvCoder = KvCoder.of( GenericRecordCoder.of(this.accountKeySchema), GenericRecordCoder.of(this.accountSchema) ); PCollection<KV<GenericRecord, GenericRecord>> kvpCollection1 = results.get(BuildGenericKafkaMessage.accountTerritoryTag) .setCoder(kvCoder);校验测试数据的Schema一致性
确保generateGenericRecord方法生成的GenericRecord,与Pipeline中实际输出的GenericRecord使用完全一致的Schema(包括字段顺序、类型、名称)。Avro对Schema匹配要求严格,微小差异都会导致编码失败。手动验证Coder的编码解码逻辑
编写小测试确认Coder能正常处理你的GenericRecord:GenericRecord expectedKey = generateGenericRecord(this.accountKeySchema, this.accountTerritory); GenericRecord expectedValue = generateGenericRecord(this.accountSchema, this.accountTerritory); // 测试key的编码解码 byte[] keyBytes = CoderUtils.encodeToByteArray(GenericRecordCoder.of(this.accountKeySchema), expectedKey); GenericRecord decodedKey = CoderUtils.decodeFromByteArray(GenericRecordCoder.of(this.accountKeySchema), keyBytes); assert expectedKey.equals(decodedKey); // 测试value的编码解码 byte[] valueBytes = CoderUtils.encodeToByteArray(GenericRecordCoder.of(this.accountSchema), expectedValue); GenericRecord decodedValue = CoderUtils.decodeFromByteArray(GenericRecordCoder.of(this.accountSchema), valueBytes); assert expectedValue.equals(decodedValue);若此步骤抛出异常,需修复
generateGenericRecord方法的Schema匹配逻辑。替换PAssert验证方式(可选)
若直接使用containsInAnyOrder仍有问题,可改用自定义ParDo手动校验结果:GenericRecord expectedKey = generateGenericRecord(this.accountKeySchema, this.accountTerritory); GenericRecord expectedValue = generateGenericRecord(this.accountSchema, this.accountTerritory); kvpCollection1.apply(ParDo.of(new DoFn<KV<GenericRecord, GenericRecord>, Void>() { @ProcessElement public void processElement(ProcessContext c) { KV<GenericRecord, GenericRecord> actual = c.element(); assertTrue(actual.getKey().equals(expectedKey)); assertTrue(actual.getValue().equals(expectedValue)); } }));
内容的提问来源于stack exchange,提问作者Prasad

