Apache Pulsar v2.0.0-rc1解密失败:如何跳过其他应用目标消息?
Pulsar消费者配置DISCARD未生效,无法跳过其他应用加密消息的解决方法
问题场景
生产者使用不同公钥对同一Topic的消息加密,供持有对应公私钥的不同应用消费。某应用消费者可正常处理自身目标消息,但消费其他应用的消息时抛出异常:
SEVERE: [topicname] [subscriptionname] Failed to decrypt data key xx.pem to decrypt messages unable to process block
Sep 04, 2022 9:57:40 AM org.apache.pulsar.client.impl.MessageCrypto decryptDataKey
消费者创建时已设置ConsumerCryptoFailureAction.DISCARD,但该配置未生效,需要实现自动跳过非自身目标的消息。
用户提供的示例代码:
try { org.apache.pulsar.client.api.Consumer consumer = pulsarClient.newConsumer() .consumerName(consumerName) .topic(topicname) .subscriptionName(subscriptionname) .subscriptionType("Shared") .cryptoKeyReader(new RawFileKeyReader(publickey , privatekey)) .cryptoFailureAction(ConsumerCryptoFailureAction.DISCARD) .subscribe(); } catch (CryptoException e) { e.printStackTrace(); } try { consumerRecord = consumer.receive(3, TimeUnit.MILLISECONDS); } catch (CryptoException e) { e.printStackTrace(); }
问题原因
ConsumerCryptoFailureAction.DISCARD的设计逻辑是当解密失败时,由Pulsar客户端自动丢弃该消息,不将其返回给应用层。但你在调用receive()方法时手动捕获了CryptoException,这会打断客户端的自动处理流程,导致解密失败的消息没有被自动跳过,反而抛出异常暴露给应用层。
解决方案
1. 移除receive()方法的CryptoException捕获
不要手动捕获解密相关的CryptoException,让客户端按照配置的DISCARD动作自动处理失败消息,这样解密失败的消息会被直接跳过,不会返回给你的业务代码。
2. 修改后的代码示例
try { org.apache.pulsar.client.api.Consumer consumer = pulsarClient.newConsumer() .consumerName(consumerName) .topic(topicname) .subscriptionName(subscriptionname) .subscriptionType(SubscriptionType.Shared) // 用枚举代替字符串,避免配置错误 .cryptoKeyReader(new RawFileKeyReader(publickey, privatekey)) .cryptoFailureAction(ConsumerCryptoFailureAction.DISCARD) .subscribe(); } catch (CryptoException e) { e.printStackTrace(); } // 仅捕获通用的Pulsar客户端异常,不处理CryptoException try { consumerRecord = consumer.receive(3, TimeUnit.MILLISECONDS); if (consumerRecord != null) { // 处理成功解密的消息 // ... 你的业务逻辑 consumer.acknowledge(consumerRecord); } } catch (PulsarClientException e) { // 处理连接超时等非解密类异常 e.printStackTrace(); }
3. 额外优化点
- 使用
SubscriptionType.Shared枚举而非字符串"Shared",避免因拼写错误导致订阅类型配置失效。 - 确保你的Pulsar客户端版本在2.8.0及以上,早期版本中
DISCARD动作的处理存在已知bug,升级到稳定版可解决这类问题。 - 检查
RawFileKeyReader加载的公私钥文件路径和权限,确保自身目标消息的解密流程正常。
内容的提问来源于stack exchange,提问作者Sekhar
相关产品推荐
相关产品推荐

