如何在Kafka Streams中删除Avro数据字段
在Kafka Streams中删除Avro数据的指定字段
你之前用filterNot的思路完全错了——filterNot是用来丢弃整条消息的,不是修改消息里的字段。你写的value == value.getEmail()这种判断逻辑永远不成立(一个对象和它的某个字段值不可能相等),所以所有消息都会被保留,字段自然不会被删除。
正确的做法是用mapValues,通过转换每条消息的内容,只保留需要的字段。下面分两种场景给出解决方案:
场景1:使用Avro Maven插件生成的具体Java类
如果已经通过插件生成了对应原Avro schema的Java类(比如叫UserOrder),最稳妥的方式是定义一个精简版的Avro schema,只包含userid、totalcost、address三个字段,再用插件生成对应的精简类(比如UserOrderSlim),然后在mapValues里做字段映射:
avrodata.peek((key, value) -> System.out.println("incoming " + value)) .mapValues(value -> { UserOrderSlim slimOrder = new UserOrderSlim(); slimOrder.setUserid(value.getUserid()); slimOrder.setTotalcost(value.getTotalcost()); slimOrder.setAddress(value.getAddress()); return slimOrder; }) .peek((key, value) -> System.out.println("processed " + value));
如果不想新增Avro类,也可以用原类的Builder模式(前提是原schema允许email和orderid为null):
avrodata.peek((key, value) -> System.out.println("incoming " + value)) .mapValues(value -> UserOrder.newBuilder() .setUserid(value.getUserid()) .setTotalcost(value.getTotalcost()) .setAddress(value.getAddress()) // 不设置email和orderid,或显式设为null .build()) .peek((key, value) -> System.out.println("processed " + value));
场景2:使用GenericRecord处理Avro数据
如果是用GenericAvroSerde处理通用Avro记录,可以直接创建新的GenericRecord,只保留需要的字段:
import org.apache.avro.Schema; import org.apache.avro.generic.GenericData; import org.apache.avro.generic.GenericRecord; // 先定义精简后的Avro Schema Schema slimSchema = new Schema.Parser().parse("{\n" + " \"type\": \"record\",\n" + " \"name\": \"UserOrderSlim\",\n" + " \"fields\": [\n" + " {\"name\": \"userid\", \"type\": \"string\"},\n" + " {\"name\": \"totalcost\", \"type\": \"double\"},\n" + " {\"name\": \"address\", \"type\": \"string\"}\n" + " ]\n" + "}"); avrodata.peek((key, value) -> System.out.println("incoming " + value)) .mapValues(value -> { GenericRecord slimRecord = new GenericData.Record(slimSchema); slimRecord.put("userid", value.get("userid")); slimRecord.put("totalcost", value.get("totalcost")); slimRecord.put("address", value.get("address")); return slimRecord; }) .peek((key, value) -> System.out.println("processed " + value));
关键总结
filterNot的作用是过滤整条消息,不能用来删除字段- 要修改消息内容、删除字段,必须用
mapValues转换每条记录 - 优先使用类型化的Avro生成类,比GenericRecord更安全,避免运行时字段名错误
内容的提问来源于stack exchange,提问作者Anis
相关产品推荐
相关产品推荐

