如何从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控制台创建导出作业:
- 进入Pinpoint控制台的分析 > 导出页面
- 点击「创建导出作业」,选择时间范围(1周前的起止日期)
- 配置导出目标为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函数完成推送:
- 在S3存储桶的属性 > 事件通知中,添加新通知,触发条件设为「所有对象创建事件」
- 选择目标为Lambda函数,编写Lambda代码读取S3中的导出文件,批量发送至Kinesis
- 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)
- Pinpoint:
内容的提问来源于stack exchange,提问作者Shahryar Faraji
相关产品推荐
相关产品推荐

