Arbeiten mit dem Firehose-Delivery-Stream

Datenstreaming mit AWS Kinesis und Lambda

Maksim Pecherskiy

Data Engineer

Bereit zum Erstellen des Streams

KF414O.png

Datenstreaming mit AWS Kinesis und Lambda

Bereit zum Erstellen des Streams

YgVbe.png

Datenstreaming mit AWS Kinesis und Lambda

Bereit zum Erstellen des Streams

BBYB4.png

Datenstreaming mit AWS Kinesis und Lambda

Bereit zum Erstellen des Streams

Ue73e.png

Datenstreaming mit AWS Kinesis und Lambda

Bereit zum Erstellen des Streams

qUq0E.png

Datenstreaming mit AWS Kinesis und Lambda

Role-ARN abrufen

qH99c.png

Datenstreaming mit AWS Kinesis und Lambda

boto3-Client initialisieren

import boto3

firehose = boto3.client('firehose', 
                        aws_access_key_id=AWS_KEY_ID, 
                        aws_secret_access_key=AWS_SECRET, 
                        region_name='us-east-1')
Datenstreaming mit AWS Kinesis und Lambda

Stream erstellen!

res = firehose.create_delivery_stream(

DeliveryStreamName = "gps-delivery-stream",
DeliveryStreamType = "DirectPut",
S3DestinationConfiguration = {
"RoleARN": "arn:aws:iam::0000000:role/firehoseDeliveryRole",
"BucketARN": "arn:aws:s3:::sd-vehicle-data"
}
)
Datenstreaming mit AWS Kinesis und Lambda

Antwort: Stream erstellt

print(res['DeliveryStreamARN'])
# ARN des neuen Streams
"arn:aws:firehose:us-east-1:0000000:deliverystream/gps-delivery-stream"
Datenstreaming mit AWS Kinesis und Lambda

Stream ist bereit

z0Fp4.png

Datenstreaming mit AWS Kinesis und Lambda

In den Stream schreiben

La0zZ.png

Datenstreaming mit AWS Kinesis und Lambda

Telematik-Hardware

OGHUHC.png

Datenstreaming mit AWS Kinesis und Lambda

Telematikdaten senden

wCJz8.png

Datenstreaming mit AWS Kinesis und Lambda

Einzelner Datensatz

{

'record_id': '939ed1d1-1740-420c-8906-445278573c7f', # <-- Eindeutige Datensatz-ID
'timestamp': '4:25:06.000', # <-- Messzeit
'vin': '4FTEX4944AK844294', # <-- Fahrzeug-ID
'lon': 106.9447146, # <-- Längengrad des Fahrzeugs
'lat': -6.3385652, # <-- Breitengrad des Fahrzeugs
'speed': 25 # <-- Fahrzeuggeschwindigkeit
}
Datenstreaming mit AWS Kinesis und Lambda

Eingehende Datensätze

vRg6A.png

Datenstreaming mit AWS Kinesis und Lambda

Weiterer Use Case

gu6oT.png

Datenstreaming mit AWS Kinesis und Lambda

Weiterer Use Case

XFD7E.png

Datenstreaming mit AWS Kinesis und Lambda

Muster

h7hUU.png

Datenstreaming mit AWS Kinesis und Lambda

Datensatz senden

res = firehose.put_record(

DeliveryStreamName='gps-delivery-stream',
Record = { 'Data': payload }
)
Datenstreaming mit AWS Kinesis und Lambda

Datensatz senden

    Record = {
        'Data': payload 
    }
Datenstreaming mit AWS Kinesis und Lambda

Datensatz senden

So sieht unser Datensatz aus
record = {
 'record_id': '939ed1d1-1740-420c-8906-445278573c7f',
 'timestamp': '4:25:06.000','vin': '4FTEX4944AK844294',
 'lon': 106.9447146,'lat': -6.338565200000001,
 'speed': 25}
Was wir senden wollen (ein String)
"939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 
4FTEX4944AK844294 106.9447146 -6.338565200000001 25"
Datenstreaming mit AWS Kinesis und Lambda

Datensatz senden

payload = " ".join(
    str(value) for value in record.values()
)
print(payload)
"939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 
4FTEX4944AK844294 106.9447146 -6.338565200000001 25"
Datenstreaming mit AWS Kinesis und Lambda

Alles zusammenführen

record = {
 'record_id': '939ed1d1-1740-420c-8906-445278573c7f',
 'timestamp': '4:25:06.000','vin': '4FTEX4944AK844294',
 'lon': 106.9447146,'lat': -6.338565200000001, 'speed': 25}

payload = " ".join( str(value) for value in record.values() )
#"939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 4FTEX4944AK844294 106.9447146 -6.338565200000001 25"
Datenstreaming mit AWS Kinesis und Lambda

Alles zusammenführen

res = firehose.put_record(
    DeliveryStreamName='gps-delivery-stream',
    Record = {
        'Data': payload + "\n" #<-- Zeilenumbruch!
    }
)
Datenstreaming mit AWS Kinesis und Lambda

Erstellte Dateien

McYNJ.png

Datenstreaming mit AWS Kinesis und Lambda

Beispieldaten

939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 4FTEX4944AK844294 106.9447146 -6.338565200000001 25
f29a5b3d-d0fa-43c0-9e1a-e2a5cdb8be7a 8:10:47.000 3FTEX1G5XAK844393 108.58068100000001 34.79925 37
ff8e7131-408d-463b-8d07-d016419b0656 20:26:44.000 2LAXX1C8XAK844292 114.39239199999999 36.097577 90
bc75da5f-1bf6-444c-80ad-49c180e1b8de 23:16:06.000 3FTEX1G5XAK844393 -76.6990172 2.481207 40
7bdcf779-444e-4313-83da-140461933aeb 22:28:44.000 5FTEX1MAXAK844295 -47.0145295 -21.4649238 40
Datenstreaming mit AWS Kinesis und Lambda

Erstellte Dateien

X0px9.png

Datenstreaming mit AWS Kinesis und Lambda

S3-Client erstellen

# Create boto3 S3 client.
s3 = boto3.client('s3', 
                  aws_access_key_id=AWS_KEY_ID, 
                  aws_secret_access_key=AWS_SECRET, 
                  region_name='us-east-1')
Datenstreaming mit AWS Kinesis und Lambda

Daten in DataFrame laden

# Get the object from S3
obj_data = s3.get_object(Bucket='sd-vehicle-data', Key=KEY_YOU_COPIED)
# Read read the object into a DataFrame
vehicle_data = pd.read_csv(
    data['Body'], 
    delimiter = " ", 
    names=["record_id", "timestamp", "vin", "lon", "lat", "speed"]))
Datenstreaming mit AWS Kinesis und Lambda

vehicle_data

record_id timestamp vin lon lat speed
0 939ed1d1... 4:25:06.000 4FTEX4944AK844294 106.945 -6.33857 25
1 f29a5b3d... 8:10:47.000 3FTEX1G5XAK844393 108.581 34.7993 37
2 ff8e7131... 20:26:44.000 2LAXX1C8XAK844292 114.392 36.0976 90
3 bc75da5f... 23:16:06.000 3FTEX1G5XAK844393 -76.699 2.48121 40
4 7bdcf779... 22:28:44.000 5FTEX1MAXAK844295 -47.0145 -21.4649 40
Datenstreaming mit AWS Kinesis und Lambda

Rückblick

qUq0E.png

Datenstreaming mit AWS Kinesis und Lambda

Rückblick

res = firehose.create_delivery_stream(

DeliveryStreamName = "gps-delivery-stream",
DeliveryStreamType = "DirectPut",
S3DestinationConfiguration = {
"RoleARN": "arn:aws:iam::0000000:role/firehoseDeliveryRole",
"BucketARN": "arn:aws:s3:::sd-vehicle-data",
}
)
Datenstreaming mit AWS Kinesis und Lambda

Rückblick

wCJz8.png

Datenstreaming mit AWS Kinesis und Lambda

Rückblick

Datenstreaming mit AWS Kinesis und Lambda

Rückblick

record_id timestamp vin lon lat speed
0 939ed1d1... 4:25:06.000 4FTEX4944AK844294 106.945 -6.33857 25
1 f29a5b3d... 8:10:47.000 3FTEX1G5XAK844393 108.581 34.7993 37
2 ff8e7131... 20:26:44.000 2LAXX1C8XAK844292 114.392 36.0976 90
3 bc75da5f... 23:16:06.000 3FTEX1G5XAK844393 -76.699 2.48121 40
4 7bdcf779... 22:28:44.000 5FTEX1MAXAK844295 -47.0145 -21.4649 40
Datenstreaming mit AWS Kinesis und Lambda

Lass uns üben!

Datenstreaming mit AWS Kinesis und Lambda

Preparing Video For Download...