消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载导读本文以 Apache Pulsar 官方 Cookbook 文档site2/website-next/versioned_docs/version-2.3.0/cookbooks-encryption.md为主线系统讲解 Pulsar 消息加密Message Encryption的实现原理与落地步骤。读完本文你将掌握Pulsar 如何用 AES ECDSA/RSA 混合加密机制保护消息载荷、如何用 OpenSSL 生成密钥对、如何实现CryptoKeyReader接口并配置加密的生产者与消费者、如何做密钥轮换与多密钥加密以及密钥丢失、批量消息等故障场景的处理策略。一、Pulsar 消息加密概述Pulsar 消息加密允许应用在生产者侧对消息加密、在消费者侧解密。加密基于应用配置的公钥/私钥对完成只有持有有效私钥的消费者才能解密消息。这一能力完全在客户端库内实现Pulsar Broker 不参与加解密只负责存储和转发密文因此 Broker 侧无法看到明文内容。关键结论来自文档与源码双重印证Pulsar 服务端任何位置都不存储加密密钥。如果私钥丢失或被删除消息将永久不可恢复加密对 Broker、BookKeeper 存储层完全透明密文以普通消息的形式持久化。二、非对称 对称混合加密原理Pulsar 采用对称加密 非对称加密的混合hybrid方案Pulsar 使用动态生成的对称 AES 密钥加密消息数据这把 AES 密钥被称为data key数据密钥data key 再用应用提供的 ECDSA/RSA 公钥/私钥对进行加密因此无需把所有消费者都共享同一把密钥生产者持有公钥公钥用于加密 AES data key消费者持有私钥私钥用于解密 data key再进而解密消息体。一次完整的加密消息生命周期如下生产者生成随机的 AES data key用公钥加密该 data key得到Encrypted AES Key加密后的 data key 作为**消息头message header/metadata**的一部分随消息发送只有持有私钥的实体即消费者才能解开 data key从而解密消息体。从源码可以看到这条链路的具体落点MessageCryptoBcpulsar-client-messagecrypto-bc/src/main/java/org/apache/pulsar/client/impl/crypto/MessageCryptoBc.java是默认的加密实现数据加密使用AES/GCM/NoPaddingAESGCM常量源码第 91 行AES 密钥长度优先取 256 位受限环境回退到 JVM 允许的最大长度data key 的非对称加密根据公钥算法分派RSA 密钥使用RSA/NONE/OAEPWithSHA1AndMGF1PaddingRSA_TRANSEC 密钥使用ECIES每条加密消息的初始化向量iv通过SecureRandom生成并写入消息元数据的encryptionParam字段加密后的 data key 以EncryptionKeys列表形式写入MessageMetadata随消息头一起传输源码encrypt()方法第 377–445 行。一个重要的设计点一条消息可以用多把密钥同时加密。生产者将同一个 data key 分别用多把公钥加密后全部写入消息头消费者只需持有其中任意一把私钥即可解密整条消息。生产者侧流程消费者侧流程图片来源site2/docs/assets/pulsar-encryption-producer.jpg、site2/docs/assets/pulsar-encryption-consumer.jpg三、快速开始从密钥生成到加密 Producer / Consumer第 1 步生成 ECDSA 公钥/私钥对使用 OpenSSL 生成 ECDSAsecp521r1 曲线密钥对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第一条命令生成私钥文件test_ecdsa_privkey.pemPEM 格式含显式椭圆曲线参数第二条命令从私钥导出 PKCS#8 格式的公钥文件test_ecdsa_pubkey.pem。文档同样支持RSA 密钥对。源码MessageCryptoBc.loadPublicKey()/loadPrivateKey()第 177–278 行通过 Bouncy Castle 的PEMParser解析 PEM 内容兼容显式 EC 参数、命名曲线 OID、X509 证书与PEMKeyPair等多种格式并处理 ECDSA 算法密钥的参数重建。第 2 步密钥管理将公钥与私钥接入你的密钥管理系统KMS / 密钥库并配置生产者检索公钥消费者客户端检索私钥。第 3 步实现 CryptoKeyReader 接口Pulsar 客户端通过CryptoKeyReader接口加载密钥。接口定义在 pulsar-client-api/src/main/java/org/apache/pulsar/client/api/CryptoKeyReader.java包含两个方法getPublicKey(String keyName, MapString, String metadata)生产者实现返回公钥getPrivateKey(String keyName, MapString, String metadata)消费者实现返回私钥。两个方法都返回EncryptionKeyInfo对象pulsar-client-api/src/main/java/org/apache/pulsar/client/api/EncryptionKeyInfo.java该对象封装了密钥字节数组setKey()以及可选的密钥元数据metadata可用于携带版本号、时间戳等额外信息。需要注意接口 Javadoc 明确提示该方法会在生产者创建时以及消费者接收消息时被调用实现中不应包含阻塞调用。第 4、5 步为 Producer / Consumer 配置加密添加加密密钥名conf.addEncryptionKey(myapp.key)注入密钥读取器实现conf.setCryptoKeyReader(keyReader)。第 6 步加密生产者示例完整可运行下面的示例完整继承自原文档可直接对照使用仓库中也提供了等价的新版 Builder 风格实现见 pulsar-client/src/test/java/org/apache/pulsar/client/tutorial/SampleCryptoProducer.javaclass RawFileKeyReader implements CryptoKeyReader { String publicKeyFile ; String privateKeyFile ; RawFileKeyReader(String pubKeyFile, String privKeyFile) { publicKeyFile pubKeyFile; privateKeyFile privKeyFile; } Override public EncryptionKeyInfo getPublicKey(String keyName, MapString, 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, MapString, 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; } } 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();第 7 步加密消费者示例完整可运行class RawFileKeyReader implements CryptoKeyReader { String publicKeyFile ; String privateKeyFile ; RawFileKeyReader(String pubKeyFile, String privKeyFile) { publicKeyFile pubKeyFile; privateKeyFile privKeyFile; } Override public EncryptionKeyInfo getPublicKey(String keyName, MapString, 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, MapString, 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; } } 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();说明ConsumerConfiguration同样只需设置CryptoKeyReader即可消费侧并不需要addEncryptionKey解密所需的密钥名来自消息头中的EncryptionKeys元数据。新版 Builder 风格 API 对照在仓库当前版本中推荐使用ProducerBuilder/ConsumerBuilder链式 API。ProducerBuilder提供了对应的.cryptoKeyReader(...)与.addEncryptionKey(...)方法见 pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerBuilder.java 第 394 行附近的addEncryptionKey示例见 SampleCryptoProducer.javaProducerbyte[] producer pulsarClient.newProducer() .topic(persistent://my-tenant/my-ns/my-topic) .cryptoKeyReader(new RawFileKeyReader(test_ecdsa_pubkey.pem, test_ecdsa_privkey.pem)) .addEncryptionKey(myappkey) .create();四、多密钥加密让多个消费者共享同一主题当消息会被跨应用边界消费时需要确保其他应用的消费者能拿到至少一把可解密的私钥。实现方式有两种消费方提供公钥消费者应用把公钥交给生产者生产者将其加入自己的加密密钥列表授权私钥生产者直接授权其密钥对中的一把私钥给消费方。对于一条消息同时用多把密钥加密的场景只需把全部密钥名都加入配置。消费者只要能访问至少一把密钥即可解密conf.addEncryptionKey(myapp.messagekey1); conf.addEncryptionKey(myapp.messagekey2);从源码看MessageCryptoBc.addPublicKeyCipher()会为每个密钥名分别用对应公钥加密同一 data key并把全部密文写入encryptedDataKeyMapConcurrentHashMapkeyName, EncryptionKeyInfoencrypt()时逐一把它们写入消息元数据的EncryptionKeys列表。消费侧getKeyAndDecryptData()第 533–559 行则遍历消息头中的所有密钥项命中缓存中的 data key 或成功解密任一项即可完成整条消息的解密。五、密钥轮换Key Rotation文档明确Pulsar 每 4 小时或发布一定数量的消息后会重新生成 AES data key生产者每 4 小时自动通过CryptoKeyReader::getPublicKey()重新拉取最新版本的非对称公钥。源码中的实现证据pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java 第 213–229 行if (this.msgCrypto ! null) { // Regenerate data key cipher at fixed interval keyGeneratorTask client.eventLoopGroup().scheduleWithFixedDelay(() - { msgCrypto.addPublicKeyCipher(conf.getEncryptionKeys(), conf.getCryptoKeyReader()); }, 0L, 4L, TimeUnit.HOURS); }即生产者创建后立即执行一次addPublicKeyCipher生成新 data key 并缓存各公钥加密结果此后每 4 小时定时刷新。与之对应MessageCryptoBc内部还维护了一个dataKeyCacheGuavaLoadingCache设置expireAfterAccess(4, TimeUnit.HOURS)用于消费者侧缓存已解密的 data key加速批量消息的解密。实践中密钥轮换配合密钥管理系统如定期更换密钥版本、由CryptoKeyReader::getPublicKey按版本返回新公钥即可实现平滑的密钥滚动。六、在生产者应用启用加密的注意事项如果生产的消息会被其他应用消费务必保证消费方应用能访问其中一把私钥见上文多密钥加密部分。在多密钥 至少一把可解的模型下任何新增消费方只需将其公钥加入生产者的addEncryptionKey列表即可在不影响既有消费者的情况下加入加密消息的解密。七、消费者应用解密加密消息消费者需要访问与消息加密公钥对应的私钥之一才能解密。如果你希望接收加密消息正确姿势是生成自己的公钥/私钥对将公钥交给生产者应用让生产者用你的公钥加密消息在自己的消费者中通过setCryptoKeyReader(...)注入能返回该私钥的CryptoKeyReader实现。八、故障处理Handling Failures8.1 生产者 / 消费者失去密钥访问权生产者侧加密操作失败时发送动作将失败并指明失败原因。应用可选择继续发送未加密消息通过conf.setCryptoFailureAction(ProducerCryptoFailureAction)控制行为。默认行为是 FAIL发送请求失败。ProducerCryptoFailureAction枚举定义在 pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ProducerCryptoFailureAction.java取值含义FAIL默认选项加密失败则发送失败SEND忽略加密失败继续以明文发送消息消费者侧因解密失败或消费者缺少密钥导致消费失败时应用可选择消费加密消息或丢弃通过conf.setCryptoFailureAction(ConsumerCryptoFailureAction)控制。默认行为是 FAIL消费请求失败。ConsumerCryptoFailureAction枚举定义在 pulsar-client-api/src/main/java/org/apache/pulsar/client/api/ConsumerCryptoFailureAction.java取值含义FAIL默认选项解密不成功则消费失败DISCARD静默确认该消息不投递给应用CONSUME将密文消息原样投递给应用由应用自行解密此时投递的加密消息带有EncryptionContext包含加密与压缩信息8.2 私钥永久丢失应用将永远无法解密这些消息。Pulsar 服务端不保存任何密钥副本因此私钥丢失等同于数据丢失——这是使用端到端加密必须接受的代价务必做好私钥的备份与安全托管。8.3 批量消息Batch Messaging的特殊性如果解密失败且消息是批量消息batch客户端将无法从批中拆出单条消息即使把ConsumerCryptoFailureAction设为CONSUME消费仍然会失败因为批量容器整体无法解析。这一点在ConsumerCryptoFailureAction.CONSUME的 Javadoc 中亦有明确说明。8.4 解密失败与积压Backlog一旦解密失败消息消费会停止此时除了客户端日志中的解密失败报错应用还会观察到积压backlog持续增长。如果应用始终无法获得解密所需的私钥唯一出路是跳过/丢弃积压中的消息例如配合DISCARD策略或清理订阅。九、进阶从源码理解默认实现与后续增强默认实现类ProducerImpl默认使用MessageCryptoBcBouncy Castle 实现由模块 pulsar-client-messagecrypto-bc 提供若运行环境缺少该模块会记录 MessageCryptoBc may not included in the jar 错误ProducerImpl.java 第 193–211 行。密钥格式MessageCryptoBc的loadPublicKey/loadPrivateKey通过 Bouncy CastlePEMParser读取 PEM 文本支持 EC含显式参数/命名曲线与 RSA密钥类型不受支持时会抛出CryptoException。IV 与密文结构每条消息的 IV 是随机生成的SecureRandom优先NativePRNGNonBlocking写入消息元数据encryptionParamdata key 密文与密钥名写入EncryptionKeys列表元数据附带的 KeyValue 可携带版本等附加信息。消费者 data key 缓存消费者按加密 data key 的 MD5 摘要为键缓存解密后的 AES data keydataKeyCache4 小时不访问即过期避免每条消息都做一次非对称解密提升吞吐。后续版本的便捷实现仓库中还提供了开箱即用的DefaultCryptoKeyReaderpulsar-client/src/main/java/org/apache/pulsar/client/impl/DefaultCryptoKeyReader.java支持通过文件路径、http(s)://、file://、data:等 URL 形式加载密钥可按密钥名分别配置公钥/私钥或提供默认密钥——在较新版本中无需手写文件读取逻辑。需要注意的是文档对应版本2.3.0的 API 仍以手写CryptoKeyReader实现为主DefaultCryptoKeyReader属于后续版本增强请以实际所用客户端版本为准。十、总结Pulsar 的消息加密通过AES 加密数据 ECDSA/RSA 加密 data key的混合方案把密钥管理完全交给应用侧Broker 全程无感知实现了真正的端到端加密。实践中的关键动作可以浓缩为四步用 OpenSSL 生成密钥对 → 实现CryptoKeyReadergetPublicKey/getPrivateKey→ 生产者addEncryptionKeysetCryptoKeyReader→ 消费者setCryptoKeyReader。在此基础上再结合多密钥加密、4 小时密钥轮换与ProducerCryptoFailureAction/ConsumerCryptoFailureAction故障策略即可在真实业务中安全、可控地落地端到端加密。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar 端到端消息加密实战指南基于 CryptoKeyReader 的公私钥加密机制Apache Pulsar 端到端消息加密实战指南基于 CryptoKeyReader 的公私钥加密机制 导读 本文围绕 Apache Pulsar 的 Pu消息队列后端流处理Apache Pulsar Node.js 客户端完整指南安装、配置与生产者/消费者/Reader 实战Apache Pulsar Node.js 客户端完整指南安装、配置与生产者/消费者/Reader 实战 Apache Pulsar 的 Node.js 客户消息队列后端流处理Apache Pulsar WebSocket API 实战指南生产者、消费者与 Reader 端点的 JSON 消息交互Apache Pulsar WebSocket API 实战指南生产者、消费者与 Reader 端点的 JSON 消息交互 本篇技术指南聚焦 Apache P消息队列后端流处理上一篇SendMIDI源码深度解析JUCE框架下的MIDI命令处理架构下一篇Django-Advanced-Filters自定义模板指南打造个性化的过滤界面创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
