- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
Pulsar 加密(Pulsar Encryption)允许应用在消息生产端对消息进行加密、在消费端进行解密,整个加解密过程由 Pulsar 客户端库在本地完成,Broker 仅负责存储与转发密文,全程接触不到明文与密钥。本文以 Apache Pulsar 2.3.2 版本官方 Cookbook 文档为主体,结合本仓库客户端源码(CryptoKeyReader.java、MessageCryptoBc.java、ProducerImpl.java、ConsumerImpl.java)与官方示例(SampleCryptoProducer.java、SampleCryptoConsumer.java),完整讲解混合加密原理、密钥生成、接口实现、Producer/Consumer 配置、密钥轮换与异常处理,帮助你快速在生产环境中落地端到端加密方案。
一、端到端加密的核心原理
1.1 为什么采用"非对称 + 对称"混合加密
Pulsar 端到端加密采用的是混合加密(Hybrid Encryption)架构:
- 对称加密(AES):Pulsar 动态生成 AES 对称密钥(即数据密钥 data key)来加密消息正文(data)。AES 加密速度快,适合加密大批量消息负载。
- 非对称加密(RSA/ECDSA):这个 AES 数据密钥本身再使用应用提供的 RSA 或 ECDSA 非对称密钥对进行加密。
采用这种两层结构的核心收益是无需把密钥分发给所有参与者:对称密钥(AES data key)只由 Producer 临时生成、用公钥加密后随消息携带,消费端只需持有对应的私钥即可恢复出 AES 密钥,中间任何人(包括 Broker)都无法解出明文。
1.2 密钥对中的角色分工
密钥(Key)指的是用于加解密的一对公钥/私钥:
- Producer 使用的公钥:来自该密钥对中的公钥部分,用于加密 AES 数据密钥;
- Consumer 使用的私钥:来自同一密钥对中的私钥部分,用于解密出 AES 数据密钥,进而解密消息。
整体流程如下:
- Producer 端应用配置公钥,Pulsar 客户端用该公钥加密 AES data key;
- 加密后的 data key 作为消息头(message header)的一部分随消息发送;
- 消息到达 Broker 后原样存储,Broker 无法解密;
- Consumer 端只有持有对应私钥的实体才能解开 data key,进而解密消息正文。
1.3 多密钥加密与消息头结构
一条消息可以同时使用多个密钥进行加密。消息的加密密钥列表(EncryptionKeys)会写入消息元数据MessageMetadata,每个条目包含密钥名称(key name)、加密后的 data key(value)以及可选的元数据(如版本、时间戳等)。
在解密端,只要持有其中任意一个用于加密该消息的私钥,即可成功解密整条消息。从源码实现看,MessageCryptoBc.decrypt() 会遍历msgMetadata.getEncryptionKeysList()中的每个加密密钥,先尝试用缓存中的 data key 解密;失败则逐个调用CryptoKeyReader.getPrivateKey()加载私钥、decryptDataKey()解出 data key,只要有一个密钥成功即可完成解密。
1.4 密钥不由 Pulsar 服务端保管
Pulsar 服务端(Broker)不会在任何地方存储加密密钥。密钥完全由应用侧通过CryptoKeyReader接口自行管理。这一点有重要的运维含义:一旦私钥丢失或被删除,消息将永久无法解密、不可恢复(irretrievably lost)。因此生产环境中务必对私钥进行备份与分级保管。
二、Producer 端加密流程(图解)
如上图所示,Producer 端加密发生在客户端本地:
- 客户端临时生成 AES 对称密钥(data key);
- 用 AES 密钥加密消息明文
Payload,得到Encrypted Payload; - 用应用配置的非对称公钥(Public Key)加密 AES 密钥,得到
Encrypted AESKey, Client key name; - 将加密后的负载与加密后的 AES 密钥组装为完整的
Encrypted message; - 加密消息经网络(Wire)发送给 Pulsar Broker 存储。
三、Consumer 端解密流程(图解)
Consumer 端解密流程与之对称:
- 加密消息从 Broker 经网络(Wire)传输到 Consumer 客户端库;
- 客户端从消息中取出公钥加密过的
Encrypted AESKey与密钥名称; - 客户端调用
CryptoKeyReader.getPrivateKey()获取私钥(Private Key),解出 AES 数据密钥; - 用 AES 密钥解密
Encrypted Payload,恢复原始明文Payload。
四、快速上手:六步启用端到端加密
官方文档给出的启用步骤如下:
- 创建 ECDSA 或 RSA 公钥/私钥对;
- 将公钥和私钥接入密钥管理,配置 Producer 检索公钥、Consumer 客户端检索私钥;
- 实现
CryptoKeyReader接口:Producer 端实现getPublicKey(),Consumer 端实现getPrivateKey(),Pulsar 客户端在需要时调用该接口加载密钥; - 为 Producer 配置添加加密密钥:
conf.addEncryptionKey("myapp.key"); - 为 Producer/Consumer 配置设置 CryptoKeyReader 实现:
conf.setCryptoKeyReader(keyReader); - 编写并运行 Producer/Consumer 示例程序。
下面逐项展开。
4.1 第一步:使用 openssl 生成 ECDSA 密钥对
官方文档推荐使用secp521r1 椭圆曲线生成 ECDSA 密钥对:
openssl ecparam -name secp521r1 -genkey -param_enc explicit -out test_ecdsa_privkey.pem openssl ec -in test_ecdsa_privkey.pem -pubout -outform pkcs8 -out test_ecdsa_pubkey.pem- 第一条命令生成包含椭圆曲线参数的 ECDSA 私钥文件
test_ecdsa_privkey.pem; - 第二条命令从私钥中导出 PKCS#8 格式的公钥文件
test_ecdsa_pubkey.pem。
注意:Pulsar 客户端要求密钥以PKCS#8 格式提供。从 ProducerBuilder.addEncryptionKey() 的 Javadoc 可以看到,应用应在回调中返回pkcs8 格式的密钥,因此
-outform pkcs8是必需的。
如果你更习惯使用 RSA,也可以用以下方式生成 2048 位 RSA 密钥对(原理一致,加密算法由客户端自动识别):
openssl genrsa -out test_rsa_privkey.pem 2048 openssl rsa -in test_rsa_privkey.pem -pubout -outform pkcs8 -out test_rsa_pubkey.pem从源码实现看,MessageCryptoBc 同时支持 RSA 与 ECDSA 两种公钥算法:RSA 密钥使用RSA/NONE/OAEPWithSHA1AndMGF1Padding变换加密 data key,ECDSA 密钥则使用 ECIES 变换(见 addPublicKeyCipher())。若传入其他不支持的密钥类型,会直接抛出CryptoException。
4.2 第二步与第三步:接入密钥管理并实现 CryptoKeyReader
密钥管理由应用自行选型(如 Vault、KMS 或本地文件系统)。本仓库官方示例采用本地 PEM 文件作为最简单的密钥存储,并通过实现CryptoKeyReader接口读取密钥。
CryptoKeyReader接口定义于 pulsar-client-api/src/main/java/org/apache/pulsar/client/api/CryptoKeyReader.java,包含两个方法:
| 方法 | 签名 | 用途 |
|---|---|---|
getPublicKey | EncryptionKeyInfo getPublicKey(String keyName, Map<String, String> metadata) | 返回与 keyName 对应的公钥(Producer 加密时调用) |
getPrivateKey | EncryptionKeyInfo getPrivateKey(String keyName, Map<String, String> metadata) | 返回与 keyName 对应的私钥(Consumer 解密时调用) |
两个方法都返回 EncryptionKeyInfo 对象,该对象封装了密钥字节(key)以及附加元数据(metadata,可携带版本、时间戳等信息,会被写入消息头的 EncryptionKeys 条目中)。
接口 Javadoc 特别强调:getPublicKey会在 Producer 创建时以及 Consumer 接收消息时被调用,应用不应在实现中做任何阻塞性调用(如不应在实现中同步访问网络 KMS)。
下面是一个从本地 PEM 文件读取密钥的完整实现(与官方示例一致):
class RawFileKeyReader implements CryptoKeyReader { String publicKeyFile = ""; String privateKeyFile = ""; RawFileKeyReader(String pubKeyFile, String privKeyFile) { publicKeyFile = pubKeyFile; privateKeyFile = privKeyFile; } @Override public EncryptionKeyInfo getPublicKey(String keyName, Map<String, String> keyMeta) { EncryptionKeyInfo keyInfo = new EncryptionKeyInfo(); try { keyInfo.setKey(Files.readAllBytes(Paths.get(publicKeyFile))); } catch (IOException e) { System.out.println("ERROR: Failed to read public key from file " + publicKeyFile); e.printStackTrace(); } return keyInfo; } @Override public EncryptionKeyInfo getPrivateKey(String keyName, Map<String, String> keyMeta) { EncryptionKeyInfo keyInfo = new EncryptionKeyInfo(); try { keyInfo.setKey(Files.readAllBytes(Paths.get(privateKeyFile))); } catch (IOException e) { System.out.println("ERROR: Failed to read private key from file " + privateKeyFile); e.printStackTrace(); } return keyInfo; } }4.3 第四步与第五步:配置 Producer 并发送加密消息
在ProducerConfiguration上完成两步关键配置:
prodConf.addEncryptionKey("myappkey"):注册用于加密消息的密钥名称;prodConf.setCryptoKeyReader(new RawFileKeyReader("test_ecdsa_pubkey.pem", "test_ecdsa_privkey.pem")):注入密钥读取器。
只有在encryptionKeys非空且cryptoKeyReader已设置时,加密才会被启用。这一判定逻辑在 ProducerConfigurationData.isEncryptionEnabled() 中实现。
完整示例:
PulsarClient pulsarClient = PulsarClient.create("http://localhost:8080"); ProducerConfiguration prodConf = new ProducerConfiguration(); prodConf.setCryptoKeyReader(new RawFileKeyReader("test_ecdsa_pubkey.pem", "test_ecdsa_privkey.pem")); prodConf.addEncryptionKey("myappkey"); Producer producer = pulsarClient.createProducer("persistent://my-tenant/my-ns/my-topic", prodConf); for (int i = 0; i < 10; i++) { producer.send("my-message".getBytes()); } pulsarClient.close();本仓库的 SampleCryptoProducer.java 提供了使用
PulsarClient.builder()+newProducer()链式 API 的等价写法,生产环境推荐使用该方式。
加密发生在哪一步?从 ProducerImpl 源码可以看到,消息在发送前会经过"压缩 → 加密"的处理管线:
- 若启用了压缩,消息先压缩后加密(若同时开启批量发送,则整个批量消息作为一个整体被加密);
- encryptMessage() 调用
msgCrypto.encrypt(conf.getEncryptionKeys(), conf.getCryptoKeyReader(), ...)完成加密; - 加密后的负载替换原 payload 后再封装进
SEND命令发送给 Broker。
4.4 Consumer 端配置与解密消息
Consumer 端只需设置CryptoKeyReader(提供私钥),不需要也不能设置加密密钥名称——解密所需的密钥名称来自消息头中携带的 EncryptionKeys 信息。
ConsumerConfiguration consConf = new ConsumerConfiguration(); consConf.setCryptoKeyReader(new RawFileKeyReader("test_ecdsa_pubkey.pem", "test_ecdsa_privkey.pem")); PulsarClient pulsarClient = PulsarClient.create("http://localhost:8080"); Consumer consumer = pulsarClient.subscribe("persistent://my-tenant/my-ns/my-topic", "my-subscriber-name", consConf); Message msg = null; for (int i = 0; i < 10; i++) { msg = consumer.receive(); // do something System.out.println("Received: " + new String(msg.getData())); } // Acknowledge the consumption of all messages at once consumer.acknowledgeCumulative(msg); pulsarClient.close();Consumer 端解密路径在 ConsumerImpl.decryptPayloadIfNeeded() 中实现:
- 检查消息元数据中
encryptionKeysCount,若为 0 说明消息未加密,直接返回原文; - 若未配置
CryptoKeyReader,则按ConsumerCryptoFailureAction处理(见下文第六节); - 调用
msgCrypto.decrypt(...)尝试解密:先检查 data key 缓存,未命中则逐个遍历消息头中的加密密钥,调用getPrivateKey()加载私钥解出 data key 后再解密正文。
五、密钥轮换(Key Rotation)
Pulsar 内置了两层密钥轮换机制:
- AES 数据密钥自动轮换:Producer 每隔4 小时(或发布一定数量的消息后)自动生成新的 AES data key。源码层面有两处印证:
- MessageCryptoBc 中
addPublicKeyCipher()每次调用都会keyGenerator.generateKey()重新生成 data key; - ProducerImpl 在 Producer 创建后启动一个固定周期为 4 小时的调度任务,反复调用
msgCrypto.addPublicKeyCipher(conf.getEncryptionKeys(), conf.getCryptoKeyReader())完成 data key 的重新生成与公钥加密。
- MessageCryptoBc 中
- 非对称公钥自动刷新:Producer 每 4 小时自动调用一次
CryptoKeyReader.getPublicKey()重新拉取公钥,以获取最新版本。因此你的密钥管理系统应保证同一 keyName 下返回的公钥是最新版本。
六、跨应用生产与多密钥加密
6.1 在 Producer 应用侧启用加密
如果你的消息会被跨应用边界消费,需要确保其他应用中的 Consumer 持有至少一把能够解密该消息的私钥。官方文档给出两种做法:
- 对方把公钥提供给你:Consumer 应用把它的公钥给你,你将这个公钥添加到自己的 Producer 密钥列表中;
- 你授权对方使用私钥:你把 Producer 密钥对中的某个私钥授予 Consumer 应用。
6.2 使用多个密钥同时加密一条消息
某些场景下 Producer 需要用多个密钥加密同一条消息(例如多个 Consumer 团队各自持有不同的私钥)。此时只需把所有密钥都addEncryptionKey进配置即可,Consumer 只要持有其中任意一个私钥就能解密:
conf.addEncryptionKey("myapp.messagekey1"); conf.addEncryptionKey("myapp.messagekey2");从源码角度解释:消息头MessageMetadata中的encryption_keys是一个可重复字段,MessageCryptoBc.encrypt() 会遍历所有 keyName,用各自公钥加密同一份 data key 并逐一写入消息头。
6.3 Consumer 端解密条件
Consumer 应用若要接收加密消息,必须持有至少一个用于加密该消息的私钥。如果 Consumer 尚未持有私钥,应自行生成一对公钥/私钥,并把公钥交给 Producer 应用用于后续消息加密。
七、故障处理(Handling Failures)
7.1 Producer/Consumer 丢失密钥访问权限
Producer 侧:当加密操作失败时(例如无法加载公钥),默认行为是发送请求失败。如果应用希望在这种情况下降级为发送明文,可以调用conf.setCryptoFailureAction(ProducerCryptoFailureAction)显式控制行为:
| 枚举值 | 含义 |
|---|---|
FAIL(默认) | 加密失败即发送失败 |
SEND | 忽略加密失败,继续发送未加密的消息 |
对应枚举定义见 ProducerCryptoFailureAction.java,默认值FAIL也体现在 ProducerConfigurationData 中。源码中 ProducerImpl.encryptMessage() 的注释也印证:除非配置显式允许发布未加密消息,否则加密失败会直接导致请求失败。
Consumer 侧:当解密失败或密钥缺失时,应用可以选择消费加密消息或丢弃它,通过conf.setCryptoFailureAction(ConsumerCryptoFailureAction)控制:
| 枚举值 | 含义 |
|---|---|
FAIL(默认) | 消费失败,消息不被投递给应用(且会记录为未确认,触发重投) |
DISCARD | 静默确认并丢弃消息,不投递给应用 |
CONSUME | 把加密消息投递给应用,由应用自行负责解密(投递的消息附带EncryptionContext,包含加密与压缩信息) |
对应枚举定义见 ConsumerCryptoFailureAction.java。在 ConsumerImpl.decryptPayloadIfNeeded() 中可以逐一看到三种分支的落点:CONSUME直接payload.retain()返回密文;DISCARD调用discardMessage(..., ValidationError.DecryptionError)丢弃;FAIL则将消息加入unAckedMessageTracker触发重投。
重要提醒:如果私钥被永久丢失,应用将永远无法解密这些消息——这是端到端加密的必然代价,请务必做好私钥备份。
7.2 批量消息(Batch Messaging)的约束
若解密失败且消息为批量消息(batch message),客户端无法从批中取出单条消息。因此即使conf.setCryptoFailureAction()设置为CONSUME,批量消息的消费依然会失败。这是批量消息整体加密设计带来的固有限制,ConsumerCryptoFailureAction.CONSUME 的 Javadoc 也明确指出了这一点。
7.3 解密失败时的观测与处置
一旦解密失败,消息消费会停止,此时可从两个现象判断:
- Backlog 增长:消费停滞导致未确认消息在 topic 上持续累积;
- 客户端日志中的解密失败信息:如
Failed to decrypt message、Message delivery failed since unable to decrypt incoming message等。
若应用没有可用的私钥来解密这些消息,唯一的处置方式是跳过/丢弃积压(backlogged)消息(例如将 Consumer 的失败动作配置为DISCARD)。
八、参考资料与源码索引
本文内容与仓库源码的对应关系如下,便于进一步深入:
- 官方文档原稿:cookbooks-encryption.md
- 加密流程示意图:Producer 侧 pulsar-encryption-producer.jpg、Consumer 侧 pulsar-encryption-consumer.jpg
- 客户端 API:
CryptoKeyReader接口 CryptoKeyReader.java、密钥载体 EncryptionKeyInfo.java、Builder 配置入口 ProducerBuilder.java - 加密实现:BC 提供方的加解密核心 MessageCryptoBc.java
- 客户端集成:Producer 加密与密钥轮换 ProducerImpl.java、Consumer 解密与失败处理 ConsumerImpl.java、配置默认值 ProducerConfigurationData.java
- 官方示例:Producer SampleCryptoProducer.java、Consumer SampleCryptoConsumer.java
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar 端到端消息加密实战指南:AES + ECDSA/RSA 混合加密体系详解
Apache Pulsar 端到端消息加密实战指南:AES + ECDSA/RSA 混合加密体系详解 Apache Pulsar 的消息加密(Message E
消息队列后端流处理Apache Pulsar 端到端加密实战指南:基于 ECDSA/RSA 密钥对的消息加密与解密(Producer/Consumer 全流程)
Apache Pulsar 端到端加密实战指南:基于 ECDSA/RSA 密钥对的消息加密与解密(Producer/Consumer 全流程) 导读 Pulsa
消息队列后端流处理Apache Pulsar 消息加密(Encryption)实战指南:基于 AES + ECDSA/RSA 混合加密的端到端消息保护方案
Apache Pulsar 消息加密(Encryption)实战指南:基于 AES + ECDSA/RSA 混合加密的端到端消息保护方案 Apache Puls
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考