【发布时间】:2017-12-08 03:46:10
【问题描述】:
我目前正在向 aws kinesis 流发送一系列 xml 消息,我一直在不同的项目中使用它,所以我非常有信心这个位有效。然后我编写了一个 lambda 来处理从 kinesis 流到 kinesis firehose 的事件:
import os
import boto3
import base64
firehose = boto3.client('firehose')
def lambda_handler(event, context):
deliveryStreamName = os.environ['FIREHOSE_STREAM_NAME']
# Send record directly to firehose
for record in event['Records']:
data = record['kinesis']['data']
response = firehose.put_record(
DeliveryStreamName=deliveryStreamName,
Record={'Data': data}
)
print(response)
我已将 kinesis 流设置为 lamdba 触发器,并将批量大小设置为 1,起始位置为 LATEST。
对于 kinesis firehose,我有以下配置:
Data transformation*: Disabled
Source record backup*: Disabled
S3 buffer size (MB)*: 10
S3 buffer interval (sec)*: 60
S3 Compression: UNCOMPRESSED
S3 Encryption: No Encryption
Status: ACTIVE
Error logging: Enabled
我发送了 162 个事件,并从 s3 读取它们,我设法得到的最多是 160 个,而且通常更少。我什至试图等待几个小时,以防重试发生奇怪的事情。
任何人都有使用 kinesis-> lamdba -> firehose 的经验,并且遇到过丢失数据的问题吗?
【问题讨论】:
-
你有没有想过这个问题?我自己也有类似的问题
标签: amazon-web-services amazon-kinesis amazon-kinesis-firehose