View a markdown version of this page

Tutorial: esegui le operazioni di base di Kinesis Data Streams utilizzando AWS CLI - Flusso di dati Amazon Kinesis

Le traduzioni sono generate tramite traduzione automatica. In caso di conflitto tra il contenuto di una traduzione e la versione originale in Inglese, quest'ultima prevarrà.

Tutorial: esegui le operazioni di base di Kinesis Data Streams utilizzando AWS CLI

Questa sezione descrive l'utilizzo di base di un flusso di dati Kinesis dalla riga di comando utilizzando la AWS CLI. Assicurati di avere familiarità con i concetti discussi in Terminologia e concetti di Amazon Kinesis Data Streams.

Nota

Dopo aver creato uno stream, il tuo account addebita costi nominali per l'utilizzo di Kinesis Data Streams perché Kinesis Data Streams non è idoneo per il piano gratuito. AWS Al termine di questo tutorial, elimina le risorse per evitare di incorrere in addebiti AWS . Per ulteriori informazioni, consulta Fase 4: pulizia.

Passaggio 1: crea uno stream

Per prima cosa devi creare un flusso e verificare che sia stato creato correttamente. Utilizza il seguente comando per creare un flusso denominato "Foo":

aws kinesis create-stream --stream-name Foo

Quindi, invia il comando seguente per verificare l'avanzamento della creazione del flusso:

aws kinesis describe-stream-summary --stream-name Foo

Dovresti ottenere un output simile a quello dell'esempio seguente:

{ "StreamDescriptionSummary": { "StreamName": "Foo", "StreamARN": "arn:aws:kinesis:us-west-2:123456789012:stream/Foo", "StreamStatus": "CREATING", "RetentionPeriodHours": 48, "StreamCreationTimestamp": 1572297168.0, "EnhancedMonitoring": [ { "ShardLevelMetrics": [] } ], "EncryptionType": "NONE", "OpenShardCount": 3, "ConsumerCount": 0 } }

In questo esempio, lo stream ha lo stato CREATING, il che significa che non è ancora pronto per l'uso. Ricontrolla dopo alcuni minuti e dovresti visualizzare un output simile a quello dell'esempio seguente:

{ "StreamDescriptionSummary": { "StreamName": "Foo", "StreamARN": "arn:aws:kinesis:us-west-2:123456789012:stream/Foo", "StreamStatus": "ACTIVE", "RetentionPeriodHours": 48, "StreamCreationTimestamp": 1572297168.0, "EnhancedMonitoring": [ { "ShardLevelMetrics": [] } ], "EncryptionType": "NONE", "OpenShardCount": 3, "ConsumerCount": 0 } }

In questo output ci sono informazioni che non ti servono per questo tutorial. Le informazioni importanti per ora sono "StreamStatus": "ACTIVE" quelle che indicano che lo stream è pronto per essere utilizzato e le informazioni sul singolo shard richiesto. Puoi inoltre verificare l'esistenza del tuo nuovo flusso utilizzando il comando list-streams, come mostrato qui:

aws kinesis list-streams

Output:

{ "StreamNames": [ "Foo" ] }

Fase 2: Inserire un record

Ora che disponi di un flusso attivo, puoi iniziare ad aggiungere dati. Per questo tutorial, utilizzerai il comando più semplice, put-record, che aggiunge un singolo record di dati contenente il testo "testdata" nel flusso:

aws kinesis put-record --stream-name Foo --partition-key 123 --data testdata

Questo comando, se utilizzato con esito positivo, genera un output simile al seguente:

{ "ShardId": "shardId-000000000000", "SequenceNumber": "49546986683135544286507457936321625675700192471156785154" }

Congratulazioni, hai appena aggiunto dati a un flusso! Nella fase successiva scoprirai come estrarre dati dal flusso.

Fase 3: Ottenere il record

GetShardIterator

Prima di poter ottenere dati dallo stream, devi procurarti lo shard iterator per lo shard che ti interessa. Un iteratore di shard definisce la posizione del flusso e lo shard da cui il consumer (il comando get-record in questo caso) effettuerà la lettura. Utilizzerai il get-shard-iterator comando come segue:

aws kinesis get-shard-iterator --shard-id shardId-000000000000 --shard-iterator-type TRIM_HORIZON --stream-name Foo

Ricorda che i comandi aws kinesis si basano su un'API del flusso di dati Kinesis, per cui se ti interessa ottenere maggiori informazioni sui parametri visualizzati, puoi consultare l'argomento della documentazione di riferimento delle API di GetShardIterator. L'esecuzione corretta produrrà un output simile al seguente esempio:

{ "ShardIterator": "AAAAAAAAAAHSywljv0zEgPX4NyKdZ5wryMzP9yALs8NeKbUjp1IxtZs1Sp+KEd9I6AJ9ZG4lNR1EMi+9Md/nHvtLyxpfhEzYvkTZ4D9DQVz/mBYWRO6OTZRKnW9gd+efGN2aHFdkH1rJl4BL9Wyrk+ghYG22D2T1Da2EyNSH1+LAbK33gQweTJADBdyMwlo5r6PqcP2dzhg=" }

La lunga stringa di caratteri apparentemente casuali è l'iteratore di shard (il tuo sarà diverso). È necessario inserire copy/paste l'iteratore shard nel comando get, mostrato di seguito. Gli iteratori shard hanno una durata valida di 300 secondi, che dovrebbe essere sufficiente per consentire all'iteratore shard di passare copy/paste al comando successivo. È necessario rimuovere tutte le nuove righe dall'iteratore shard prima di incollarle al comando successivo. Se ricevete un messaggio di errore che indica che l'iteratore shard non è più valido, eseguite nuovamente il comando. get-shard-iterator

GetRecords

Il comando get-records estrae i dati dal flusso e si risolve in una chiamata a GetRecords nell'API del flusso di dati Kinesis. L'iteratore di shard specifica la posizione nello shard da cui desideri iniziare a leggere i record di dati in sequenza. Se non sono disponibili record nella porzione dello shard a cui l'iteratore punta, GetRecords restituisce un elenco vuoto. Potrebbero essere necessarie più chiamate per accedere a una parte dello shard che contiene i record.

Nel seguente esempio del get-records comando:

aws kinesis get-records --shard-iterator AAAAAAAAAAHSywljv0zEgPX4NyKdZ5wryMzP9yALs8NeKbUjp1IxtZs1Sp+KEd9I6AJ9ZG4lNR1EMi+9Md/nHvtLyxpfhEzYvkTZ4D9DQVz/mBYWRO6OTZRKnW9gd+efGN2aHFdkH1rJl4BL9Wyrk+ghYG22D2T1Da2EyNSH1+LAbK33gQweTJADBdyMwlo5r6PqcP2dzhg=

Se stai eseguendo questo tutorial da un processore di Unix-type comandi come bash, puoi automatizzare l'acquisizione dello shard iterator usando un comando annidato, come questo:

SHARD_ITERATOR=$(aws kinesis get-shard-iterator --shard-id shardId-000000000000 --shard-iterator-type TRIM_HORIZON --stream-name Foo --query 'ShardIterator') aws kinesis get-records --shard-iterator $SHARD_ITERATOR

Se stai eseguendo questo tutorial da un sistema che supporta PowerShell, puoi automatizzare l'acquisizione dello shard iterator usando un comando come questo:

aws kinesis get-records --shard-iterator ((aws kinesis get-shard-iterator --shard-id shardId-000000000000 --shard-iterator-type TRIM_HORIZON --stream-name Foo).split('"')[4])

Il risultato positivo del get-records comando richiederà i record dal tuo stream per lo shard che hai specificato quando hai ottenuto lo shard iterator, come nell'esempio seguente:

{ "Records":[ { "Data":"dGVzdGRhdGE=", "PartitionKey":"123”, "ApproximateArrivalTimestamp": 1.441215410867E9, "SequenceNumber":"49544985256907370027570885864065577703022652638596431874" } ], "MillisBehindLatest":24000, "NextShardIterator":"AAAAAAAAAAEDOW3ugseWPE4503kqN1yN1UaodY8unE0sYslMUmC6lX9hlig5+t4RtZM0/tALfiI4QGjunVgJvQsjxjh2aLyxaAaPr+LaoENQ7eVs4EdYXgKyThTZGPcca2fVXYJWL3yafv9dsDwsYVedI66dbMZFC8rPMWc797zxQkv4pSKvPOZvrUIudb8UkH3VMzx58Is=" }

Nota che get-records è descritta sopra come una richiesta, il che significa che potresti ricevere zero o più record anche se ci sono record nel tuo stream. Qualsiasi record restituito potrebbe non rappresentare tutti i record attualmente presenti nel tuo stream. Questo è normale e il codice di produzione interrogherà lo stream alla ricerca di record a intervalli appropriati. Questa velocità di polling varierà in base ai requisiti di progettazione specifici dell'applicazione.

Nella registrazione di questa parte del tutorial, noterai che i dati sembrano inutili e non è il testo testdata in chiaro che abbiamo inviato. Ciò è dovuto al modo in cui put-record utilizza la codifica Base64 per consentirti di inviare dati binari. Tuttavia, il supporto di Kinesis Data Streams AWS CLI non fornisce la decodifica Base64 perché la decodifica Base64 su contenuti binari non elaborati stampati su stdout può causare comportamenti indesiderati e potenziali problemi di sicurezza su determinate piattaforme e terminali. Se utilizzate un decodificatore Base64 (ad esempio https://www.base64decode.org/) per decodificare manualmente, vedrete che in effetti è così. dGVzdGRhdGE= testdata Questo è sufficiente per il bene di questo tutorial perché, in pratica, AWS CLI viene usato raramente per consumare dati. Più spesso, viene utilizzato per monitorare lo stato dello streaming e ottenere informazioni, come mostrato in precedenza (describe-streamelist-streams). Per ulteriori informazioni sulla KCL, consulta Sviluppo di consumer personalizzati con velocità di trasmissione effettiva condivisa tramite KCL.

get-recordsnon restituisce sempre tutti i record stream/shard specificati. Se ciò si verifica, utilizza NextShardIterator dall'ultimo risultato per ottenere il set di record successivo. Se si inserissero più dati nello stream, situazione normale nelle applicazioni di produzione, è possibile continuare a eseguire il polling dei dati get-records ogni volta. Tuttavia, se non si chiama get-records utilizzando il next shard iterator entro la durata di 300 secondi dell'iteratore shard, verrà visualizzato un messaggio di errore e sarà necessario utilizzare il get-shard-iterator comando per ottenere un nuovo iteratore shard.

L'output include anche MillisBehindLatest, che è il numero di millisecondi a cui la risposta dell'operazione GetRecords si trova rispetto all'estremità del flusso, a indicare il ritardo rispetto all'ora corrente del consumer. Un valore di zero indica che l'elaborazione dei record è aggiornata e che non sono presenti nuovi record da elaborare in questo momento. Nel nostro caso, potresti riscontrare un numero piuttosto elevato se hai dedicato diverso tempo ad acquisire familiarità con le fasi di questo tutorial. Per impostazione predefinita, i record di dati rimangono in uno stream per 24 ore in attesa che tu li recuperi. Questo intervallo di tempo viene chiamato periodo di conservazione ed è configurabile fino a 365 giorni.

Un get-records risultato positivo sarà sempre NextShardIterator pari se non ci sono più record nello stream. Si tratta di un modello di polling che presume che un producer stia potenzialmente immettendo più record nel flusso in un dato momento. Sebbene possa scrivere le tue routine di polling, se utilizzi la KCL menzionata in precedenza per lo sviluppo di applicazioni consumer, il polling verrà gestito per tuo conto.

Se chiami get-records finché non ci sono più record nello stream e nello shard da cui stai prelevando, vedrai un output con record vuoti simile al seguente esempio:

{ "Records": [], "NextShardIterator": "AAAAAAAAAAGCJ5jzQNjmdhO6B/YDIDE56jmZmrmMA/r1WjoHXC/kPJXc1rckt3TFL55dENfe5meNgdkyCRpUPGzJpMgYHaJ53C3nCAjQ6s7ZupjXeJGoUFs5oCuFwhP+Wul/EhyNeSs5DYXLSSC5XCapmCAYGFjYER69QSdQjxMmBPE/hiybFDi5qtkT6/PsZNz6kFoqtDk=" }

Fase 4: pulizia

Elimina lo stream per liberare risorse ed evitare addebiti involontari sul tuo account. Esegui questa operazione ogni volta che crei uno stream e non lo utilizzerai, perché gli addebiti vengono addebitati per stream indipendentemente dal fatto che tu stia inserendo e ricevendo dati con esso o meno. Il comando clean-up è il seguente:

aws kinesis delete-stream --stream-name Foo

Il successo non comporta alcun risultato. Usa describe-stream per controllare lo stato di avanzamento dell'eliminazione:

aws kinesis describe-stream-summary --stream-name Foo

Se esegui questo comando subito dopo il comando di eliminazione, vedrai un risultato simile al seguente esempio:

{ "StreamDescriptionSummary": { "StreamName": "samplestream", "StreamARN": "arn:aws:kinesis:us-west-2:123456789012:stream/samplestream", "StreamStatus": "ACTIVE",

Quando il flusso è stato completamente eliminato, describe-stream genererà un errore di tipo "non trovato":

A client error (ResourceNotFoundException) occurred when calling the DescribeStreamSummary operation: Stream Foo under account 123456789012 not found.