
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载本文是一份围绕 Apache Pulsar 端到端加密End-to-End Encryption的实战指南基于本仓库 2.2.1 版本文档与对应客户端源码整理而成。Pulsar 端到端加密允许应用在生产端加密消息、在消费端解密消息Broker 全程只负责存储和转发密文不参与任何加解密运算。读完本文你将掌握其对称密钥加密数据、非对称密钥加密数据密钥的双层加密原理能够在 Java、C、Python、Node.js 四种客户端中完成密钥生成、CryptoKeyReader配置、多密钥加密与失败处理。端到端加密解决的问题在默认场景下Pulsar 只保证传输链路安全TLS与 Broker 侧的鉴权授权但消息在 Broker 上是以明文形式存储和转发的。对于跨组织、跨应用边界的数据交换场景或者对数据主权有强要求的业务应用希望只有持有合法密钥的消费端能够还原消息内容此时就需要端到端加密。Pulsar 端到端加密的核心约定是消息在Producer 客户端内完成加密密文随普通消息一同进入 Pulsar 服务消息在Consumer/Reader 客户端内完成解密应用拿到的是明文Pulsar 服务端不保存任何加密密钥——Broker、BookKeeper 存储层都只接触密文。如果私钥丢失或删除对应消息将永久无法解密、不可恢复。这是端到端加密与服务端托管密钥方案的本质区别也是选择该方案前必须接受的运维约束。双层加密原理对称与非对称结合Pulsar 采用混合加密Hybrid Encryption策略综合了对称加密的高性能与非对称加密的密钥分发便利性对称层加密数据Pulsar 客户端动态生成一个随机的AES 对称密钥Data Key用它对消息体Payload进行对称加密。对称加密速度快适合承载真实业务数据。非对称层加密密钥应用提供ECDSA 或 RSA 密钥对客户端用其中公钥加密 AES Data Key得到 Encrypted Data Key。组装Encrypted Data Key 作为消息头Message Header的一部分与密文一起发送。只有持有对应私钥的实体即消费端才能解出 Data Key进而解密消息体。从源码实现看这套逻辑落在MessageCrypto接口及其 Bouncy Castle 实现中。接口定义了encrypt、decrypt、addPublicKeyCipher、removeKeyCipher等方法见 MessageCrypto.java而 MessageCryptoBc.java 是实际加解密引擎从中可以确认如下实现细节数据加密算法AES/GCM/NoPaddingAES-GCM认证加密IV 长度为 12 字节GCM tag 长度为 128 bit数据密钥加密算法按公钥/私钥算法自动选择——RSA 使用RSA/NONE/OAEPWithSHA1AndMGF1PaddingEC 曲线密钥使用ECIESAES 密钥强度优先使用 256 bit若 JVM 受 JCE 无限制强度策略文件限制而getMaxAllowedKeyLength(AES) 128则退化为 128 bit 并打印告警日志解密缓存消费端对解密出的 AES Data Key 做 MD5 摘要后放入缓存expireAfterAccess(4, TimeUnit.HOURS)避免每条消息都执行一次非对称解密。密钥对的分工也很明确Producer 持有公钥用于加密 Data KeyConsumer 持有私钥用于解密 Data Key。一个生产端还可以用多个公钥同时加密同一条消息消息中会携带多份 Encrypted Data Key任意一个对应私钥在手即可完成解密这为跨应用密钥分发提供了弹性。Producer 端加密流程Pulsar Producer 端消息加密流程示意图Producer 客户端库在发送前完成生成 AES Data Key → 用公钥加密 Data Key 得到 Encrypted Data Key → 用 Data Key 加密消息体 → 组装出包含密文 Encrypted Data Key 密钥名称的加密消息再经传输链路交给 Broker。Consumer 端解密流程Pulsar Consumer 端消息解密流程示意图Consumer 客户端库收到加密消息后先用私钥解密消息头中的 Encrypted Data Key还原 AES Data Key再用它解密消息体最终把明文交付给应用。快速开始五步启用端到端加密第 1 步生成 ECDSA 或 RSA 密钥对使用 OpenSSL 生成 ECDSAsecp521r1曲线密钥对openssl ecparam -name secp521r1 -genkey -param_enc explicit -out test_ecdsa_privkey.pem openssl ec -in test_ecdsa_privkey.pem -pubout -outform pem -out test_ecdsa_pubkey.pem如需改用 RSA可用openssl genrsa生成私钥、openssl rsa -pubout导出公钥。最终产出两个 PEM 文件公钥交给 Producer 侧私钥保留在 Consumer 侧。第 2 步将密钥接入密钥管理把公私钥交给应用的密钥管理设施文件系统、KMS、Vault 等并约定Producer 应用从中获取公钥Consumer 客户端从中获取私钥。第 3 步实现 CryptoKeyReader 接口CryptoKeyReader是 Pulsar 客户端访问密钥库的统一抽象。Java 侧接口定义见 CryptoKeyReader.java它要求实现两个方法getPublicKey(String keyName, MapString, String metadata)Producer 创建与消息加密时被调用返回公钥的EncryptionKeyInfogetPrivateKey(String keyName, MapString, String metadata)Consumer 收到消息后解密时被调用返回私钥的EncryptionKeyInfo。接口注释中有两条重要提示该方法在Producer 创建时以及 Consumer 接收消息时都会被调用因此实现中不应包含阻塞调用。第 4 步为 Producer 指定加密密钥名在 Producer 构建器上通过addEncryptionKey(myapp.key)声明要使用的密钥名。密钥名只是一个逻辑标识真正的密钥字节由CryptoKeyReader按此名称加载。第 5 步将 CryptoKeyReader 配置到 Producer / Consumer / Reader不同语言客户端的配置方式如下四种语言均以persistent://my-tenant/my-ns/my-topic为例JavaPulsarClient pulsarClient PulsarClient.builder().serviceUrl(pulsar://localhost:6650).build(); String topic persistent://my-tenant/my-ns/my-topic; // RawFileKeyReader 是示例实现并非 Pulsar 官方提供 CryptoKeyReader keyReader new RawFileKeyReader(test_ecdsa_pubkey.pem, test_ecdsa_privkey.pem); Producerbyte[] producer pulsarClient.newProducer() .topic(topic) .cryptoKeyReader(keyReader) .addEncryptionKey(myappkey) .create(); Consumerbyte[] consumer pulsarClient.newConsumer() .topic(topic) .subscriptionName(my-subscriber-name) .cryptoKeyReader(keyReader) .subscribe(); Readerbyte[] reader pulsarClient.newReader() .topic(topic) .startMessageId(MessageId.earliest) .cryptoKeyReader(keyReader) .create();CC 客户端除了可自定义实现外还内置了DefaultCryptoKeyReader直接传入公私钥文件路径即可定义见 CryptoKeyReader.hsince 2.8.0Client client(pulsar://localhost:6650); std::string topic persistent://my-tenant/my-ns/my-topic; // DefaultCryptoKeyReader 是内置实现从文件读取公钥和私钥 auto keyReader std::make_sharedDefaultCryptoKeyReader(test_ecdsa_pubkey.pem, test_ecdsa_privkey.pem); Producer producer; ProducerConfiguration producerConf; producerConf.setCryptoKeyReader(keyReader); producerConf.addEncryptionKey(myappkey); client.createProducer(topic, producerConf, producer); Consumer consumer; ConsumerConfiguration consumerConf; consumerConf.setCryptoKeyReader(keyReader); client.subscribe(topic, my-subscriber-name, consumerConf, consumer); Reader reader; ReaderConfiguration readerConf; readerConf.setCryptoKeyReader(keyReader); client.createReader(topic, MessageId::earliest(), readerConf, reader);Pythonfrom pulsar import Client, CryptoKeyReader client Client(pulsar://localhost:6650) topic persistent://my-tenant/my-ns/my-topic # CryptoKeyReader 是内置实现从文件读取公钥和私钥 key_reader CryptoKeyReader(test_ecdsa_pubkey.pem, test_ecdsa_privkey.pem) producer client.create_producer( topictopic, encryption_keymyappkey, crypto_key_readerkey_reader ) consumer client.subscribe( topictopic, subscription_namemy-subscriber-name, crypto_key_readerkey_reader ) reader client.create_reader( topictopic, start_message_idMessageId.earliest, crypto_key_readerkey_reader ) client.close()Node.jsconst Pulsar require(pulsar-client); (async () { // 创建客户端 const client new Pulsar.Client({ serviceUrl: pulsar://localhost:6650, operationTimeoutSeconds: 30, }); // 创建生产者 const producer await client.createProducer({ topic: persistent://public/default/my-topic, sendTimeoutMs: 30000, batchingEnabled: true, publicKeyPath: public-key.client-rsa.pem, encryptionKey: encryption-key }); // 创建消费者 const consumer await client.subscribe({ topic: persistent://public/default/my-topic, subscription: sub1, subscriptionType: Shared, ackTimeoutMs: 10000, privateKeyPath: private-key.client-rsa.pem }); // 发送消息 for (let i 0; i 10; i 1) { const msg my-message-${i}; producer.send({ data: Buffer.from(msg), }); console.log(Sent message: ${msg}); } await producer.flush(); // 接收消息 for (let i 0; i 10; i 1) { const msg await consumer.receive(); console.log(msg.getData().toString()); consumer.acknowledge(msg); } await consumer.close(); await producer.close(); await client.close(); })();注意Node.js 客户端中 Producer 通过publicKeyPathencryptionKey指定公钥与密钥名Consumer 通过privateKeyPath指定私钥文件路径底层同样依赖CryptoKeyReader语义。自定义 CryptoKeyReader 实现当密钥存放于文件之外的存储如数据库、KMS、密钥管理平台时需要自行实现CryptoKeyReader。四种语言的现状如下。Java完整示例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; } }接入时把RawFileKeyReader实例传给.cryptoKeyReader(keyReader)即可。C骨架class CustomCryptoKeyReader : public CryptoKeyReader { public: Result getPublicKey(const std::string keyName, std::mapstd::string, std::string metadata, EncryptionKeyInfo encKeyInfo) const override { // TODO: return ResultOk; } Result getPrivateKey(const std::string keyName, std::mapstd::string, std::string metadata, EncryptionKeyInfo encKeyInfo) const override { // TODO: return ResultOk; } }; auto keyReader std::make_sharedCustomCryptoKeyReader(/* ... */); // TODO: create producer, consumer or reader based on keyReader hereC 中除自定义实现外也可以使用内置的DefaultCryptoKeyReader只需在构造时指定私钥、公钥文件路径。Python 与 Node.js目前 Python 与 Node.js 客户端不支持自定义CryptoKeyReader实现只能使用默认实现通过指定私钥、公钥文件路径来启用加密即上文快速开始中的用法。密钥轮换机制Pulsar 的密钥轮换分为两个层面AES Data Key 的轮换Producer 端每4 小时或在发布一定数量消息后生成一个新的随机 AES Data Key避免长期使用同一把对称密钥非对称公钥的轮换Producer 每隔4 小时调用一次CryptoKeyReader.getPublicKey()重新拉取公钥从而感知外部密钥管理中的最新公钥版本。在实现上MessageCryptoBc的dataKeyCache同样以expireAfterAccess(4, TimeUnit.HOURS)组织缓存消费端缓存的 Data Key 到期后会自动重新走一次私钥解 Data Key流程与服务端 4 小时的轮换周期天然对齐。这意味着只要密钥管理端及时更新公钥版本无需重启 Producer 即可完成滚动轮换。多密钥加密跨应用边界的安全协作如果 Producer 生产的消息会被多个应用的消费者读取就必须保证每个消费应用都能解密。有两种协作方式对方给你公钥消费者所在应用把它们的公钥提供给 ProducerProducer 将其加入自己的加密密钥列表你给对方私钥你从自己使用的密钥对中向对方授权其中一个私钥的访问权。当 Producer 需要同时使用多把公钥加密时通过多次调用addEncryptionKey声明所有密钥名PulsarClient.newProducer().addEncryptionKey(myapp.messagekey1).addEncryptionKey(myapp.messagekey2);此时每条消息的消息头会携带分别用myapp.messagekey1、myapp.messagekey2公钥加密的两份 Data Key。消费端只要持有其中任意一把对应的私钥即可解密该消息——这一特性从MessageCrypto.decrypt的实现也可得到印证解密时会遍历消息头中的加密密钥列表逐一尝试用本地私钥还原 Data Key命中任意一个即继续解密消息体。消费端解密要求Consumer 要成功解密消息必须能访问到 Producer 使用的至少一把公钥对应的私钥。如果你是需要接收加密消息的一方正确姿势是创建自己的公私钥对把你的公钥交给 Producer 应用Producer 用你的公钥加密消息你在 Consumer 客户端通过CryptoKeyReader提供自己的私钥完成解密。失败处理与最佳实践Producer / Consumer 丢失密钥访问权Producer 侧加密失败时发送操作会以失败告终并指明失败原因。若业务允许降级可通过PulsarClient.newProducer().cryptoFailureAction(ProducerCryptoFailureAction)控制行为。枚举定义见 ProducerCryptoFailureAction.javaFAIL默认加密失败即发送失败SEND忽略加密失败以明文继续发送消息。Consumer 侧因解密失败或缺少密钥而无法消费时可通过PulsarClient.newConsumer().cryptoFailureAction(ConsumerCryptoFailureAction)控制行为。枚举定义见 ConsumerCryptoFailureAction.javaFAIL默认持续失败直到解密成功DISCARD静默确认消息不投递给应用CONSUME把密文原样投递给应用由应用自行解密消息自带EncryptionContext包含加密与压缩信息。需要特别强调的是若私钥永久丢失无论设置哪种失败策略应用都无法再解密这些消息。批处理Batch消息的限制如果解密失败的消息是批量消息Batch Message客户端将无法从中取出单条消息因此即使把cryptoFailureAction()设置为ConsumerCryptoFailureAction.CONSUME整批消息的消费依然会失败。这是使用批量发送Batching时需要评估的边界条件。积压Backlog与排障当解密持续失败时消息消费会停滞表现为客户端日志持续打印解密失败信息同时该分区的backlog 不断增长。如果应用确实没有对应私钥唯一的选择是跳过或丢弃这些积压消息例如配合DISCARD策略或按消息 ID 跳过而不是无休止地重试。总结Pulsar 端到端加密把加密数据与加密密钥两个环节分离动态 AES Data Key 保证数据加密性能RSA/ECDSA 非对称密钥对保证密钥只在应用侧流转Broker 全程只接触密文、不持有任何密钥。落地时只需抓住三个关键点用 OpenSSL 生成密钥对、实现或使用默认的CryptoKeyReader、在 Producer/Consumer/Reader 构建器上声明密钥名与失败策略。多密钥加密、4 小时密钥轮换、失败动作枚举等机制让该方案能够覆盖跨应用协作、密钥滚动更新和故障降级等真实生产场景。相关接口与实现可在本仓库的 CryptoKeyReader.java、MessageCrypto.java 与 MessageCryptoBc.java 中继续深入研读。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar 端到端消息加密实战指南AES ECDSA/RSA 混合加密体系详解Apache Pulsar 端到端消息加密实战指南AES ECDSA/RSA 混合加密体系详解 Apache Pulsar 的消息加密Message E消息队列后端流处理Apache Pulsar 端到端加密实战指南AES ECDSA/RSA 混合加密原理与多语言实现Apache Pulsar 端到端加密实战指南AES ECDSA/RSA 混合加密原理与多语言实现 本文围绕 Apache Pulsar 的端到端加密E消息队列后端流处理终极移动端加密完全指南对称加密、非对称加密与密钥管理最佳实践终极移动端加密完全指南对称加密、非对称加密与密钥管理最佳实践 移动端应用开发中数据安全是用户信任的基石。GitHub 加速计划 / an / android文档教程移动开发上一篇PyTorch-VAE文档生成使用Sphinx构建API文档下一篇FSPagerView持续集成配置GitHub Actions自动化构建与测试创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考