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

调用PutLogEvents API时遇到nextSequenceToken无效问题求助

解决CloudWatch PutLogEvents的InvalidSequenceTokenException问题

问题背景

处理组织级CloudTrail日志时,通过Kinesis Stream订阅,尝试为每个账号创建对应日志组,调用AWS CloudWatch API的put_log_events时遇到InvalidSequenceTokenException错误。

原Python代码

import base64
import gzip
import json
import logging
import os

import boto3

# Setup logging configuration
logging.basicConfig()
logger = logging.getLogger()
logger.setLevel(logging.INFO)

logs_client = boto3.client('logs', region_name=os.getenv('AWS_REGION'))

# setting environment variables
global seq_token

def unpack_kinesis_stream_records(event):
    # decode and decompress each base64 encoded data element
    return [gzip.decompress(base64.b64decode(k["kinesis"]["data"])).decode('utf-8') for k in event["Records"]]


def decode_raw_cloud_trail_events(cloudTrailEventDataList):
    # Convert Raw Event Data List
    eventList = [json.loads(e) for e in cloudTrailEventDataList]

    # Filter out-non DATA_MESSAGES since we only require cloud watch message type = DATA_MESSAGE
    filteredEvents = [
        e for e in eventList if e["messageType"] == 'DATA_MESSAGE']

    # Covert each individual log Event Message
    events = []
    for f in filteredEvents:
        for e in f["logEvents"]:
            events.append(
                {
                    'timestamp': e["timestamp"],
                    'message': e["message"],
                }
            )

    events.sort(key=lambda x: x["timestamp"])
    logger.info("{0} Event Logs Decoded".format(len(events)))

    log_group = ("log_group_for_cloudtrail_"+eventList[0]["logStream"].split("_")[1])

    log_stream = os.getenv('AWS_LAMBDA_LOG_STREAM_NAME')

    # creating log group
    try:
        logs_client.create_log_group(
            logGroupName=log_group,
            tags={
                'Created By': 'ZH Lamda'
            }
        )
    except logs_client.exceptions.ResourceAlreadyExistsException:
        print("log group exists")

    try:
        logs_client.create_log_stream(
            logGroupName=log_group,
            logStreamName=log_stream,
        )
    except logs_client.exceptions.ResourceAlreadyExistsException:
        print("log stream already exists")

    return [log_group, log_stream, events]


def handle_request(event, context):
    seq_token = None
    # Unpack Kinesis Stream Records
    kinesis_data = unpack_kinesis_stream_records(event)
    # Decode and filter events
    events = decode_raw_cloud_trail_events(kinesis_data)

    if len(events[2]) == 0:
        return f'No events to process'

    log_event = {
        'logGroupName': events[0],
        'logStreamName': events[1],
        'logEvents': events[2],
    }

    if seq_token is not None:
        log_event['sequenceToken'] = seq_token

    response = logs_client.put_log_events(**log_event)
    seq_token = response['nextSequenceToken']

    return f"Successfully processed {len(events)} records."


def lambda_handler(event, context):
    return handle_request(event, context)

错误日志

[ERROR] InvalidSequenceTokenException: An error occurred (InvalidSequenceTokenException) when calling the PutLogEvents operation: The given sequenceToken is invalid. The next expected sequenceToken is: null
Traceback (most recent call last):
  File "/var/task/lambda_function.py", line 120, in lambda_handler
    return handle_request(event, context)
  File "/var/task/lambda_function.py", line 107, in handle_request
    response = logs_client.put_log_events(**log_event)
  File "/var/runtime/botocore/client.py", line 391, in _api_call
    return self._make_api_call(operation_name, kwargs)
  File "/var/runtime/botocore/client.py", line 719, in _make_api_call
    raise error_class(parsed_response, operation_name)

问题根源

  1. Lambda无状态导致令牌丢失:Lambda每次调用都是独立执行,seq_token变量每次都会被重置为None,但如果目标日志流已有写入记录,必须使用上次返回的nextSequenceToken才能继续写入。
  2. 日志流命名不合理:当前使用AWS_LAMBDA_LOG_STREAM_NAME作为目标日志流名称,这是Lambda自身的日志流,会导致不同批次的CloudTrail日志写入同一个流,增加令牌管理复杂度,也不符合按账号隔离日志的需求。
  3. 无异常重试逻辑:遇到令牌失效时,没有尝试获取正确令牌并重试请求。

修复方案

1. 持久化序列令牌

Lambda是无状态的,需要用外部存储(比如DynamoDB)保存每个日志流对应的nextSequenceToken,键可以用logGroupName#logStreamName。

添加DynamoDB操作代码:

import boto3
dynamodb = boto3.resource('dynamodb')
token_table = dynamodb.Table(os.getenv('TOKEN_TABLE_NAME')) # 需提前创建表,主键为`stream_key`

def get_sequence_token(log_group, log_stream):
    stream_key = f"{log_group}#{log_stream}"
    try:
        response = token_table.get_item(Key={'stream_key': stream_key})
        return response['Item']['sequence_token'] if 'Item' in response else None
    except Exception as e:
        logger.error(f"获取序列令牌失败: {e}")
        return None

def update_sequence_token(log_group, log_stream, token):
    stream_key = f"{log_group}#{log_stream}"
    try:
        token_table.put_item(Item={'stream_key': stream_key, 'sequence_token': token})
    except Exception as e:
        logger.error(f"更新序列令牌失败: {e}")

2. 修正日志流命名

替换Lambda自身日志流名称,改用基于账号和时间的命名规则,确保日志流隔离:

# 在decode_raw_cloud_trail_events函数中替换log_stream生成逻辑
import time
account_id = eventList[0]["logStream"].split("_")[1]
log_stream = f"cloudtrail_{account_id}_{int(time.time())}"

3. 添加异常重试逻辑

遇到令牌错误时,自动获取正确令牌并重试:

import time

def put_logs_with_retry(logs_client, log_event):
    max_retries = 3
    retry_count = 0
    while retry_count < max_retries:
        try:
            response = logs_client.put_log_events(**log_event)
            return response['nextSequenceToken']
        except logs_client.exceptions.InvalidSequenceTokenException as e:
            logger.info(f"令牌无效,获取正确令牌: {e}")
            # 从日志流描述中获取正确令牌
            stream_desc = logs_client.describe_log_streams(
                logGroupName=log_event['logGroupName'],
                logStreamNamePrefix=log_event['logStreamName'],
                limit=1
            )
            correct_token = stream_desc['logStreams'][0].get('uploadSequenceToken')
            log_event['sequenceToken'] = correct_token
            retry_count += 1
            time.sleep(1)
        except Exception as e:
            logger.error(f"写入日志失败: {e}")
            raise
    raise Exception("写入日志重试次数已达上限")

4. 整合后的核心处理逻辑

修改handle_request函数:

def handle_request(event, context):
    # 解析Kinesis记录
    kinesis_data = unpack_kinesis_stream_records(event)
    # 解码CloudTrail事件
    log_group, log_stream, events = decode_raw_cloud_trail_events(kinesis_data)

    if len(events) == 0:
        return '无事件需要处理'

    log_event = {
        'logGroupName': log_group,
        'logStreamName': log_stream,
        'logEvents': events,
    }

    # 获取保存的序列令牌
    seq_token = get_sequence_token(log_group, log_stream)
    if seq_token is not None:
        log_event['sequenceToken'] = seq_token

    # 带重试的日志写入
    try:
        next_token = put_logs_with_retry(logs_client, log_event)
        # 更新令牌到DynamoDB
        update_sequence_token(log_group, log_stream, next_token)
        return f"成功处理{len(events)}条记录"
    except Exception as e:
        logger.error(f"处理失败: {e}")
        return f"处理记录失败: {str(e)}"

额外注意事项

  • 确保Lambda的IAM角色拥有DynamoDB的读写权限。
  • 日志事件的timestamp必须在过去14天内且不超过当前时间2小时,否则会被CloudWatch拒绝。
  • 单次put_log_events请求最多支持10000条事件,总大小不超过1MB,事件过多时需要分批处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 05:06:28