rdkafka中FutureRecord为何需指定key/类型?key乱码如何解决
问题成因
两个问题的根源分别是:
- 编译报错:
FutureRecord<'a, K, P>是带两个泛型参数的类型,K为key的类型、P为消息体的类型。仅调用.payload()时只能确定泛型P的类型,K没有任何可用于推断的上下文,而Rust要求async块内的所有类型必须在编译期可确定,因此触发E0698类型推断错误,和你是否真的需要传key无关。 - Key乱码:你调用
to_le_bytes()得到的是i64值的原始小端二进制字节数组,不是UTF-8编码的可读文本。绝大多数Kafka消息查看工具会默认把key的字节按UTF-8字符串解析,原始整数字节不符合UTF-8编码规则,也不属于可打印ASCII字符范围,自然会显示为乱码。你在本地打印i64::from_le_bytes(offset)能得到正确值,是因为你手动把字节按整数规则反解了,和工具的解析逻辑不一样。
正确实现方案
根据你对key的使用需求选对应写法即可:
需要key为可读文本(绝大多数场景推荐)
直接将offset转为字符串,取字符串的UTF-8字节作为key即可,这样不管是消费端还是可视化工具查看都不会出现乱码:async { loop { let msg = consumer.recv().await.unwrap(); println!("Offset :{:?}", msg.offset()); let dest_topic = "compression_data_control"; let offset = msg.offset(); let offset_str = offset.to_string(); let record = FutureRecord::to(dest_topic) .payload(msg.payload().unwrap()) .key(offset_str.as_bytes()); let send_res = producer.send(record, Timeout::After(Duration::from_millis(15000))).await; send_res.map_or(-1, |_| { println!("Processed offset {} successfully.", offset); 0 }); } }.await;确实需要用二进制格式存储i64类型的key
这种场景下显示乱码是正常现象,不需要修改生产端代码,只要在消费key的逻辑中,拿到原始字节数组后调用i64::from_le_bytes()反序列化为整数即可,不要把二进制字节当UTF-8字符串解析。不需要设置key
只要显式给编译器确定K的类型即可,rust-rdkafka为单元类型()实现了ToBytes特征,对应空key的语义,写法如下:let record = FutureRecord::to(dest_topic) .payload(msg.payload().unwrap()) .key(()); // 传单元类型,明确K为(),对应无key
内容的提问来源于stack exchange,提问作者rtviii
相关产品推荐
相关产品推荐

