求AWS Node.js Lambda实现Kafka消费者解密Confluent信封加密消息的代码示例
AWS Node.js Lambda 消费Confluent信封加密Kafka消息代码示例
前置依赖
先安装所需第三方包:
npm install kafkajs @aws-sdk/client-secrets-manager
建议将所有敏感配置存入AWS Secrets Manager,硬编码密钥会有安全风险。需要配置的参数如下:
- Kafka集群broker地址列表
- 消费组ID、目标topic名称
- Kafka集群认证凭据(SASL用户名/密码,或AWS MSK IAM角色权限)
- 用于解密AES数据密钥的RSA私钥(注:信封加密流程中公钥用于加密数据密钥,私钥用于解密,若对方仅提供公钥请先确认加密流程是否符合预期)
完整实现代码
const { Kafka } = require('kafkajs'); const crypto = require('crypto'); const { SecretsManagerClient, GetSecretValueCommand } = require("@aws-sdk/client-secrets-manager"); // 全局初始化消费者,避免Lambda冷启动重复初始化 let kafkaConsumer = null; let rsaPrivateKey = null; const SECRET_NAME = process.env.SECRET_NAME; const AWS_REGION = process.env.AWS_REGION || 'us-east-1'; // 加载敏感配置 const loadSecrets = async () => { if (rsaPrivateKey) return rsaPrivateKey; const client = new SecretsManagerClient({ region: AWS_REGION }); const response = await client.send(new GetSecretValueCommand({ SecretId: SECRET_NAME })); const secrets = JSON.parse(response.SecretString); rsaPrivateKey = secrets.RSA_PRIVATE_KEY; return rsaPrivateKey; }; // 1. 初始化Kafka消费者连接 const initKafkaConsumer = async () => { if (kafkaConsumer) return kafkaConsumer; const privateKey = await loadSecrets(); const kafka = new Kafka({ clientId: 'lambda-kafka-consumer', brokers: process.env.KAFKA_BROKERS.split(','), // 如果是Confluent云或SASL认证,打开以下配置 // sasl: { // mechanism: 'plain', // username: process.env.KAFKA_USER, // password: process.env.KAFKA_PASSWORD // }, // ssl: true }); kafkaConsumer = kafka.consumer({ groupId: process.env.KAFKA_GROUP_ID }); await kafkaConsumer.connect(); await kafkaConsumer.subscribe({ topic: process.env.KAFKA_TOPIC, fromBeginning: false }); return kafkaConsumer; }; // 2. Confluent信封加密消息解密逻辑 const decryptConfluentMessage = (encryptedPayload, rsaPrivateKey) => { const payloadBuffer = Buffer.from(encryptedPayload, 'base64'); let offset = 0; // 跳过1字节版本号 offset += 1; // 读取2字节加密AES密钥长度 const aesKeyLength = payloadBuffer.readUInt16BE(offset); offset += 2; // 读取加密的AES密钥 const encryptedAesKey = payloadBuffer.slice(offset, offset + aesKeyLength); offset += aesKeyLength; // 读取16字节AES初始化向量IV const iv = payloadBuffer.slice(offset, offset + 16); offset += 16; // 读取实际加密的业务数据 const encryptedData = payloadBuffer.slice(offset); // 用RSA私钥解密得到AES 256密钥 const aesKey = crypto.privateDecrypt( { key: rsaPrivateKey, padding: crypto.constants.RSA_PKCS1_OAEP_PADDING, oaepHash: "sha256", }, encryptedAesKey ); // 用AES密钥解密业务数据 const decipher = crypto.createDecipheriv('aes-256-cbc', aesKey, iv); let decryptedData = decipher.update(encryptedData); decryptedData = Buffer.concat([decryptedData, decipher.final()]); return JSON.parse(decryptedData.toString('utf8')); }; // Lambda入口函数 exports.handler = async (event, context) => { // 允许Lambda等待异步操作完成再退出 context.callbackWaitsForEmptyEventLoop = true; const consumer = await initKafkaConsumer(); const privateKey = await loadSecrets(); const processedRecords = []; // 3. 轮询消费Kafka记录 await consumer.run({ eachMessage: async ({ topic, partition, message }) => { try { const decryptedData = decryptConfluentMessage(message.value, privateKey); processedRecords.push({ topic, partition, offset: message.offset, data: decryptedData }); // 这里添加你的业务处理逻辑 console.log('处理成功:', decryptedData); } catch (err) { console.error('消息处理失败:', err, { offset: message.offset, partition }); throw err; } }, // 单次轮询最大处理记录数,可根据Lambda内存配置调整 eachBatchAutoResolve: true, partitionsConsumedConcurrently: 3 }); // 如果是事件源映射触发的Lambda,不需要手动轮询,直接处理event中的记录即可,解密逻辑通用 // 以下为事件源映射触发时的处理示例: // for (const record of event.records) { // const decryptedData = decryptConfluentMessage(record.value, privateKey); // 业务处理逻辑 // } return { statusCode: 200, body: JSON.stringify({ processedCount: processedRecords.length }) }; };
注意事项
- Lambda执行角色需要配置对应权限:Secrets Manager读权限、Kafka集群访问权限(如果是AWS MSK需要添加
kafka:DescribeCluster,kafka:GetBootstrapBrokers,kafka:Consume等权限) - 若使用Lambda Kafka事件源映射,可删除手动初始化消费者和轮询的逻辑,直接处理event传入的记录即可,解密逻辑可以复用
- 如果生产者使用的加密填充方式、IV长度和示例不一致,请根据实际情况调整解密参数
内容的提问来源于stack exchange,提问作者intelviji
相关产品推荐
相关产品推荐

