基于KSTREAMS Java实现按薪资过滤Kafka主题数据
使用Kafka Streams过滤薪资高于30000的记录
你已经走完了最关键的前两步,接下来的过滤逻辑其实很直接,我会给你完整的代码示例和细节注意点,帮你快速完成开发。
先衔接你的现有步骤
假设你的输入记录格式是类似 员工ID,姓名,薪资,部门 这样的逗号分隔文本(比如 1,张三,35000,技术部),你的前两步代码大概是这样:
StreamsBuilder builder = new StreamsBuilder(); // 步骤1:读取输入主题并分割文本 KStream<String, String> inputStream = builder.stream("你的输入主题名称"); KStream<String, String[]> splitStream = inputStream.mapValues(record -> record.split(",")); // 步骤2:将薪资字段设为Key // 注意:这里假设薪资是分割后的第3个元素(索引为2,从0开始计数),请根据实际格式调整索引 KStream<String, String[]> keyedBySalary = splitStream.selectKey((originalKey, splitRecord) -> splitRecord[2]);
步骤3:实现过滤逻辑并输出到目标主题
接下来我们要做的就是筛选薪资>30000的记录,然后转回原始格式发送到输出主题,同时处理可能的异常(比如薪资不是数字的情况):
// 步骤3-1:过滤薪资高于30000的记录 KStream<String, String[]> filteredStream = keyedBySalary.filter((salaryStr, splitRecord) -> { try { // 把薪资字符串转成数值类型(整数用Integer,带小数用Double) double salary = Double.parseDouble(salaryStr); return salary > 30000; } catch (NumberFormatException e) { // 处理薪资格式错误的记录:可以选择丢弃,或者转发到专门的错误主题 System.err.println("发现无效薪资格式:" + salaryStr + ",记录内容:" + String.join(",", splitRecord)); return false; // 返回false表示丢弃这条记录 } }); // 步骤3-2:将分割后的数组转回逗号分隔的原始文本格式 KStream<String, String> outputStream = filteredStream.mapValues(splitRecord -> String.join(",", splitRecord)); // 步骤3-3:发送到输出主题 outputStream.to("你的输出主题名称");
几个关键注意点
- 字段索引确认:一定要根据你实际的记录格式调整薪资字段的索引(比如如果格式是
姓名,薪资,部门,那索引就是1),否则会取错值导致过滤逻辑失效。 - 异常处理:实际生产环境中,建议把格式错误的记录转发到专门的错误主题,而不是只打印日志,这样方便后续排查和处理脏数据。
- 数值类型选择:如果薪资是整数,用
Integer.parseInt()更高效;如果有小数(比如30000.5),就用Double.parseDouble()。 - Kafka Streams基础配置:别忘了在启动流应用时配置必要参数,示例如下:
Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "salary-filter-app"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass()); KafkaStreams streams = new KafkaStreams(builder.build(), props); streams.start(); // 优雅关闭钩子 Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
这样整个流程就完整了,启动应用后,输入主题中薪资高于30000的记录会自动被过滤到输出主题里。
内容的提问来源于stack exchange,提问作者Outlander
相关产品推荐
相关产品推荐

