如何在Spark中读取Kafka主题的二进制加密消息并防数据损坏
如何在Spark Streaming中读取Kafka的二进制加密数据
你目前用的是Spark Streaming基于Receiver的Kafka读取API,默认用字符串解码器处理消息,这肯定会破坏二进制加密数据——毕竟字符串编码转换会篡改原始字节。只需要把解码器换成字节数组解码器,就能完美保留原始加密数据,修改方式很简单:
1. 导入字节数组解码器类
首先得引入Kafka的字节数组解码器:
import org.apache.spark.streaming.kafka.ByteArrayDecoder;
2. 修改KafkaUtils.createStream调用
把原来的字符串泛型替换为byte[],并指定解码器参数:
// 替换原有的String,String泛型为byte[],byte[],同时传入ByteArrayDecoder JavaPairReceiverInputDStream<byte[], byte[]> pairrdd = KafkaUtils.createStream( jssc, zkQuorum, group, topicmap, ByteArrayDecoder.class, ByteArrayDecoder.class );
关键说明
ByteArrayDecoder会直接读取Kafka消息的原始二进制字节,不会做任何字符串编码转换,彻底避免了加密数据因编码转换导致的损坏问题- 之后你可以直接操作
pairrdd中的byte[]类型value(也就是加密的二进制数据),进行解密等后续处理 - 如果后续想切换到Spark Streaming的Direct API(更推荐的生产级方式),思路完全一致:把解码器换成
ByteArrayDecoder,泛型改为byte[]即可
额外注意
确保你的项目依赖中包含了对应版本的Spark-Kafka集成包,比如针对Spark 2.x、Scala 2.11的依赖:spark-streaming-kafka-0-8_2.11(版本号要和你的Spark版本匹配)
内容的提问来源于stack exchange,提问作者kevin
相关产品推荐
相关产品推荐

