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

如何从AWS Pinpoint检索1周前事件并发送至Kinesis数据流

从AWS Pinpoint检索历史事件并发送至Kinesis Data Stream

确认事件可用性

  • AWS Pinpoint默认将事件数据保留90天,1周前的事件大概率仍在存储中。你可以在Pinpoint控制台的项目设置 > 数据管理页面,查看当前的事件保留配置,确认数据未被清理。

检索1周前的事件

方式1:使用GetEventsBatch API(小批量数据)

适合按端点/用户ID定向获取事件,可通过AWS CLI或SDK调用:

# AWS CLI 示例,替换参数为你的实际信息
aws pinpoint get-events-batch \
  --application-id "your-pinpoint-app-id" \
  --request-body '{
    "EndpointIds": ["endpoint-id-1", "endpoint-id-2"],
    "StartDateTime": "2024-05-01T00:00:00Z",
    "EndDateTime": "2024-05-07T23:59:59Z"
  }'
  • StartDateTime和EndDateTime需设为1周前的时间范围,格式遵循ISO 8601标准。

方式2:导出全量历史事件(大批量数据)

如果需要获取指定时间段的全量事件,可通过Pinpoint控制台创建导出作业:

  1. 进入Pinpoint控制台的分析 > 导出页面
  2. 点击「创建导出作业」,选择时间范围(1周前的起止日期)
  3. 配置导出目标为S3存储桶,完成后等待导出完成

将事件发送至Kinesis Data Stream

方式1:API获取后直接推送(小批量)

用Python SDK编写脚本,将GetEventsBatch获取到的事件推送到Kinesis:

import boto3
import json

# 初始化客户端
pinpoint_client = boto3.client('pinpoint')
kinesis_client = boto3.client('kinesis')

# 获取1周前的事件
response = pinpoint_client.get_events_batch(
    ApplicationId='your-pinpoint-app-id',
    RequestBody={
        'EndpointIds': ['target-endpoint-id'],
        'StartDateTime': '2024-05-01T00:00:00Z',
        'EndDateTime': '2024-05-07T23:59:59Z'
    }
)

# 批量推送至Kinesis
for endpoint_id, event_data in response['EventsResponse']['Results'].items():
    for event in event_data['Events'].values():
        kinesis_client.put_record(
            StreamName='your-kinesis-stream-name',
            Data=json.dumps(event).encode('utf-8'),
            PartitionKey=endpoint_id
        )

方式2:S3导出后触发推送(大批量)

当事件导出到S3后,通过S3事件通知触发Lambda函数完成推送:

  1. 在S3存储桶的属性 > 事件通知中,添加新通知,触发条件设为「所有对象创建事件」
  2. 选择目标为Lambda函数,编写Lambda代码读取S3中的导出文件,批量发送至Kinesis
  3. Lambda示例代码核心逻辑:
    import boto3
    import json
    import gzip
    
    s3 = boto3.client('s3')
    kinesis = boto3.client('kinesis')
    
    def lambda_handler(event, context):
        bucket = event['Records'][0]['s3']['bucket']['name']
        key = event['Records'][0]['s3']['object']['key']
        
        # 读取S3中的压缩导出文件
        response = s3.get_object(Bucket=bucket, Key=key)
        with gzip.GzipFile(fileobj=response['Body']) as f:
            events = json.load(f)
            
        # 批量推送至Kinesis
        records = [{'Data': json.dumps(evt).encode('utf-8'), 'PartitionKey': evt['EndpointId']} for evt in events]
        kinesis.put_records(Records=records, StreamName='your-kinesis-stream-name')
    

权限配置注意事项

  • 确保执行操作的IAM角色/用户拥有以下权限:
    • Pinpoint:mobileanalytics:GetEvents、pinpoint:CreateExportJob
    • Kinesis:kinesis:PutRecord、kinesis:PutRecords
    • S3:s3:GetObject、s3:PutObject(若使用导出到S3)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 15:30:58