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

如何通过Lambda函数向MSK Topic发送消息

实现Lambda向AWS MSK集群Topic发送消息的完整方案

一、补充Terraform配置

1. 输出MSK集群的Bootstrap地址

在你的aws_msk_cluster资源后添加输出,让Lambda能获取集群连接地址:

output "msk_bootstrap_brokers" {
  value = aws_msk_cluster.this.bootstrap_brokers
}

2. 完善Lambda环境变量

将MSK的Bootstrap地址传入Lambda环境变量,修改aws_lambda_function的environment块:

environment {
  variables = {
    KAFKA_TOPIC       = var.kafka_topics[each.key]
    BOOTSTRAP_SERVERS = aws_msk_cluster.this.bootstrap_brokers
    # 若启用IAM认证,需添加集群ARN
    # MSK_CLUSTER_ARN  = aws_msk_cluster.this.arn
  }
}

3. 配置安全组规则

确保Lambda与MSK网络连通:

  • Lambda安全组允许出站到MSK的Kafka端口(明文用9092,TLS用9094)
  • MSK安全组允许入站来自Lambda安全组的对应端口流量

添加Terraform安全组规则资源:

# MSK允许Lambda入站访问Kafka端口
resource "aws_security_group_rule" "msk_lambda_inbound" {
  type                     = "ingress"
  from_port                = 9092 # 若用TLS则改为9094
  to_port                  = 9092
  protocol                 = "tcp"
  security_group_id        = var.security_groups_id[0] # MSK的安全组ID
  source_security_group_id = var.sg[0] # Lambda的安全组ID
}

4. 精简Lambda IAM权限

替换你当前宽泛的权限,保留最小必要权限:

  • AWSLambdaBasicExecutionRole:用于CloudWatch日志输出
  • AWSLambdaVPCAccessExecutionRole:用于VPC内资源访问
  • 若启用IAM认证,添加AmazonMSKFullAccess(或自定义细粒度策略)

二、Lambda Python代码实现消息发送

1. 依赖打包准备

Lambda默认不包含Kafka客户端库,需本地打包依赖:

# 创建匹配Python 3.8的虚拟环境
python3.8 -m venv venv
source venv/bin/activate
# 安装兼容MSK 2.8.2的kafka-python版本
pip install kafka-python==2.0.2
# 复制依赖到当前目录
cp -r venv/lib/python3.8/site-packages/* .
# 打包代码+依赖为zip
zip -r lambda_deployment.zip lambda_function.py .

将生成的lambda_deployment.zip作为Terraform中filename的值。

2. 核心消息发送代码

修改你的Python脚本,实现Kafka消息生产逻辑:

import json
import logging
import os
from kafka import KafkaProducer
from kafka.errors import KafkaError

# 初始化日志
logger = logging.getLogger()
logger.setLevel(logging.INFO)

def lambda_handler(event, context):
    # 从环境变量读取配置
    bootstrap_servers = os.environ['BOOTSTRAP_SERVERS'].split(',')
    kafka_topic = os.environ['KAFKA_TOPIC']
    
    logger.info(f"Received event: {json.dumps(event)}")
    logger.info(f"Target: Bootstrap servers={bootstrap_servers}, Topic={kafka_topic}")

    try:
        # 初始化Kafka生产者(明文连接示例)
        producer = KafkaProducer(
            bootstrap_servers=bootstrap_servers,
            value_serializer=lambda v: json.dumps(v).encode('utf-8'),
            # 若使用TLS加密,取消以下注释
            # security_protocol='SSL',
            # ssl_check_hostname=False
        )

        # 处理事件数据(假设event包含待发送的data字段)
        message_data = event.get('data', {})
        # 发送消息并等待确认
        future = producer.send(kafka_topic, value=message_data)
        record_metadata = future.get(timeout=10)
        
        logger.info(f"Message sent: Topic={record_metadata.topic}, Partition={record_metadata.partition}, Offset={record_metadata.offset}")
        
        producer.flush()
        producer.close()

        return {
            'statusCode': 200,
            'message': 'Message sent successfully'
        }
    except KafkaError as e:
        logger.error(f"Kafka error occurred: {str(e)}")
        return {'statusCode': 500, 'message': f"Kafka error: {str(e)}"}
    except Exception as e:
        logger.error(f"Unexpected error: {str(e)}")
        return {'statusCode': 500, 'message': f"Error: {str(e)}"}

三、可选:启用MSK IAM认证(更安全)

若需避免明文连接,配置IAM认证:

1. 修改MSK集群配置

添加SASL IAM认证:

resource "aws_msk_cluster" "this" {
  # 原有配置...
  client_authentication {
    sasl {
      iam {
        enabled = true
      }
    }
  }
}

2. 更新Lambda代码

安装IAM认证依赖kafka-aws-iam-sasl-signer,并修改生产者配置:

# 新增导入
import boto3
from kafka_aws_iam_sasl_signer import IAMClient

# 初始化IAM签名客户端
cluster_arn = os.environ['MSK_CLUSTER_ARN']
signer = IAMClient(cluster_arn)

# 替换生产者初始化代码
producer = KafkaProducer(
    bootstrap_servers=bootstrap_servers,
    security_protocol='SASL_SSL',
    sasl_mechanism='AWS_MSK_IAM',
    sasl_plain_username='',
    sasl_plain_password=signer.get_token(),
    value_serializer=lambda v: json.dumps(v).encode('utf-8')
)

四、验证步骤

  1. 部署Terraform配置,确认MSK集群和Lambda创建成功
  2. 触发Lambda(可通过测试事件传入{"data": "test message"})
  3. 查看CloudWatch日志,确认消息发送成功
  4. 可选:使用Kafka消费者工具验证Topic中存在目标消息

内容的提问来源于stack exchange,提问作者Homer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 11:58:10