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

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)
报错根因
  1. 类型传参错误:pymysql默认游标返回的查询结果是元组列表,遍历得到的单行row是tuple类型,而base64.b64encode()仅支持字节类对象入参,直接传入元组会直接触发该类型错误。
  2. 逻辑冗余错误:Kinesis Firehose的写入接口不需要手动做base64编码,boto3 SDK会自动处理传输层编码;对编码结果额外调用json.dumps()属于多余操作,会导致数据格式不符合预期。
  3. 调用逻辑不合理:循环内每次仅传1条记录却调用批量写入接口put_record_batch,既浪费接口配额,也容易触发限流,且没有做失败记录校验,存在数据丢失风险。
  4. 连接配置隐患:原代码把RDS连接定义在handler全局作用域,Lambda执行环境复用时闲置连接会自动断开,后续请求会出现连接失效问题。
修复方案
  1. 游标替换为DictCursor,让查询返回的单行结果为带字段名的字典,方便下游解析数据结构。
  2. 先将单行字典序列化为JSON字符串,编码为UTF-8字节后直接传入Data字段,不需要手动做base64编码。
  3. 收集所有待写入记录后统一调用批量写入接口,单批数据不超过Firehose的接口限制(单批最多500条、总大小不超过4MiB)。
  4. 数据库连接放到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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 20:12:21