You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

求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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 21:15:02