Python实现AWS Lambda向Kinesis Firehose推送RDS数据报TypeError
问题描述
基于Lambda实现RDS到Kinesis Firehose的数据推送时,pymysql连接RDS查询成功,但写入Firehose阶段触发类型错误:
- 错误类型:
TypeError - 错误信息:
"a bytes-like object is required, not 'tuple'"
原问题代码如下:
connection = pymysql.connect(host = endpoint, user = username, passwd = password, db = database_name) FIREHOSE_STREAM = 'DEMOLAMBDAFIREHOSE' client = boto3.client('firehose') def lambda_handler(event, context): cursor = connection.cursor() cursor.execute('SELECT * from inventory.report_product') rows = cursor.fetchall() for row in rows: data = base64.b64encode(row) response = client.put_record_batch( DeliveryStreamName=FIREHOSE_STREAM, Records=[ { 'Data': json.dumps(data) }, ] ) print (response)
报错根因
- 类型传参错误:pymysql默认游标返回的查询结果是元组列表,遍历得到的单行
row是tuple类型,而base64.b64encode()仅支持字节类对象入参,直接传入元组会直接触发该类型错误。 - 逻辑冗余错误:Kinesis Firehose的写入接口不需要手动做base64编码,boto3 SDK会自动处理传输层编码;对编码结果额外调用
json.dumps()属于多余操作,会导致数据格式不符合预期。 - 调用逻辑不合理:循环内每次仅传1条记录却调用批量写入接口
put_record_batch,既浪费接口配额,也容易触发限流,且没有做失败记录校验,存在数据丢失风险。 - 连接配置隐患:原代码把RDS连接定义在handler全局作用域,Lambda执行环境复用时闲置连接会自动断开,后续请求会出现连接失效问题。
修复方案
- 游标替换为
DictCursor,让查询返回的单行结果为带字段名的字典,方便下游解析数据结构。 - 先将单行字典序列化为JSON字符串,编码为UTF-8字节后直接传入
Data字段,不需要手动做base64编码。 - 收集所有待写入记录后统一调用批量写入接口,单批数据不超过Firehose的接口限制(单批最多500条、总大小不超过4MiB)。
- 数据库连接放到handler内初始化,请求结束后主动关闭,避免失效连接报错。
修复后的可运行代码:
import json import boto3 import pymysql from pymysql.cursors import DictCursor # 配置项 endpoint = "你的RDS实例访问端点" username = "数据库账号" password = "数据库密码" database_name = "目标库名" FIREHOSE_STREAM = 'DEMOLAMBDAFIREHOSE' # 全局初始化Firehose客户端(Lambda中客户端全局初始化可复用,提升性能) firehose_client = boto3.client('firehose') def lambda_handler(event, context): # 数据库连接放到handler内初始化,避免复用失效连接 conn = pymysql.connect( host=endpoint, user=username, passwd=password, db=database_name, cursorclass=DictCursor ) try: with conn.cursor() as cursor: cursor.execute('SELECT * from inventory.report_product') rows = cursor.fetchall() # 组装Firehose写入记录 records = [] for row in rows: # 序列化后加换行符,方便Firehose投递到下游(如S3)后按行解析 row_json = json.dumps(row, ensure_ascii=False) + "\n" records.append({ "Data": row_json.encode("utf-8") }) # 批量写入,数据量超过500条时需要做切分处理 if records: resp = firehose_client.put_record_batch( DeliveryStreamName=FIREHOSE_STREAM, Records=records ) # 可按需添加失败记录重试逻辑 print(f"写入完成,失败记录数:{resp.get('FailedPutCount', 0)}") finally: # 确保连接关闭 conn.close()
额外注意事项
- 如果查询返回数据量较大,需要按单批最多500条、总大小不超过4MiB的规则切分批次写入,否则会触发Firehose接口参数错误。
- 如果业务场景确实需要对数据做base64编码,必须先将序列化后的JSON字符串转为UTF-8字节,再传入
base64.b64encode(),禁止直接传入元组、字典这类结构化对象。 - 生产环境建议添加失败记录重试逻辑,通过接口返回的
FailedPutCount和RequestResponses字段定位写入失败的记录,重新投递避免数据丢失。
内容的提问来源于stack exchange,提问作者Darshan Pratheep
相关产品推荐
相关产品推荐

