【发布时间】:2020-01-19 14:13:02
【问题描述】:
我正在尝试将数据发送到流,然后使用 kinesis firehouse 将数据传递到 ElasticSearch,我正在使用 python lambda 函数在推送之前将数据转换为 JSON,但是 lambda 失败并出现以下错误。
[ERROR] KeyError: 'Records'
Traceback (most recent call last):
File "/var/task/lambda_function.py", line 7, in lambda_handler
for record in event["Records"]:
我可以使用下面的 sharditerator 查看分片中的记录。
{
"Records": [
{
"SequenceNumber": "49599580114447666780699883212202628138922281244234350610",
"ApproximateArrivalTimestamp": 1568741427.415,
"Data": "MjAwNi8wMS8wMSAwMDowMDowMHwzMTA4IE9DQ0lERU5UQUwgRFJ8M3wzQyAgICAgICAgfDExMTV8MTA4NTEoQSlWQyBUQUtFIFZFSCBXL08gT1dORVJ8MjQwNHwzOC41NTA0MjA0N3wtMTIxLjM5MTQxNTh8MjAxOS8wOS8xNyAyMzowMDoyNA==",
"PartitionKey": "1"
},
我正在使用下面的 lambda 函数来处理流。
import json
print("Loading the function")
success = 0
failure = 0
def lambda_handler(event, context):
for record in event["Records"]:
print(record)
payload=base64.b64decode(record["Data"]).decode('utf-8')
match = payload.split('|')
result = {}
if match:
# create a dict object of the row
#build all fields from array
result["crime_time"] = match[0]
result["address"] = match[1]
result['district'] = int(match[2])
result['beat'] = match[3]
result['grid'] =int(match[4])
result['description'] = match[5]
result['crime_id'] = int(match[6])
result['latitude'] = float(match[7])
result['longitude'] = float(match[8])
result['load_time'] = match[9]
result['location'] = {
'lat' : float(match[7]),
'lon' : float(match[8])
}
success+=1
return {
'statusCode': 200,
'body': json.dumps(result)
}
但是在将数据发送到流后,我在 lambda 函数中遇到错误。
【问题讨论】:
标签: python-3.x amazon-web-services aws-lambda amazon-kinesis-firehose