View a markdown version of this page

Amazon Managed Service per Apache Flink 2.2 - Servizio gestito per Apache Flink

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à.

Amazon Managed Service per Apache Flink 2.2

Amazon Managed Service per Apache Flink ora supporta Apache Flink versione 2.2. Questo è il primo importante aggiornamento della versione del servizio. Questa pagina illustra le funzionalità introdotte in Flink 2.2, insieme a importanti considerazioni per l'aggiornamento da Flink 1.x.

Nota

Flink 2.2 introduce modifiche sostanziali che richiedono un'attenta pianificazione. Consulta l'elenco completo delle principali modifiche e deprecazioni qui sotto e Guida alla compatibilità dello stato per gli aggiornamenti di Flink 2.2 prima di effettuare l'aggiornamento dalla versione 1.x.

Amazon Managed Service for Apache Flink 2.2 introduce modifiche comportamentali che potrebbero interrompere le applicazioni esistenti al momento dell'aggiornamento. Esaminali attentamente insieme alle modifiche all'API Flink nella sezione successiva.

Gestione programmatica della configurazione

Rimozione delle metriche

  • La fullRestarts metrica è stata rimossa in Flink 2.2. Usa invece la numRestarts metrica.

  • La bytesRequestedPerFetch metrica per il connettore KDS è stata rimossa nella versione 6.0.0 del AWS connettore Flink (solo la versione del connettore compatibile con Flink 2.2).

  • Le downtime metriche uptime e sono entrambe contrassegnate come obsolete in Flink 2.2 e verranno rimosse presto. Sostituisci uptime con la nuova metrica. runningTime Sostituisci downtime con uno o più tra restartingTimecancellingTime, efailingTime.

  • Vedi la pagina Metriche e dimensioni per l'elenco completo delle metriche supportate.

Non-Credential Chiamate IMDS bloccate

  • Questi endpoint consentiti vengono utilizzati dagli AWS SDK DefaultCredentialsProvider (/latest/meta-data/iam/security-credentials/) e DefaultAwsRegionProviderChain (/latest/dynamic/instance-identity/document) per configurare automaticamente le credenziali e la regione per l'applicazione.

  • Le applicazioni che utilizzano le funzioni dell'SDK AWS che si basano su chiamate IMDS senza credenziali (comeEC2MetadataUtils.getInstanceId(), EC2MetadataUtils.getInstanceType()EC2MetadataUtils.getLocalHostName(), oEC2MetadataUtils.getAvailabilityZone()) riceveranno errori HTTP 4xx durante il tentativo di queste chiamate.

  • Se l'applicazione utilizza IMDS per esempio metadati o altre informazioni al di fuori dei percorsi consentiti, esegui il refactoring del codice per utilizzare invece le variabili di ambiente o la configurazione dell'applicazione.

Read-Only File system principale

  • Per migliorare la sicurezza, qualsiasi dipendenza al di fuori della /tmp quale è la directory di lavoro predefinita di flink risulterà in:. java.io.FileNotFoundException: /{path}/{filename} (Read-only file system)

  • Le dipendenze del file system possono provenire direttamente dal codice o indirettamente dalle librerie incluse nelle dipendenze. Sostituisci le dipendenze dirette del file system nel tuo codice. /tmp/ Per le dipendenze indirette del file system dalle librerie, utilizza le sostituzioni della configurazione della libreria per reindirizzare le operazioni del file system a. /tmp/

Di seguito è riportato un riepilogo delle principali modifiche e deprecazioni introdotte in Managed Service for Apache Flink 2.2. Consulta le note di rilascio di Apache Flink 2.0 per le note di rilascio complete di Apache Flink 2.0 che introducono queste importanti modifiche.

Rimozioni di API e linguaggi di Flink

DataSet API rimossa

  • L' DataSet API precedente per l'elaborazione in batch è stata completamente rimossa in Flink 2.0+. Tutte le elaborazioni in batch devono ora utilizzare l'API DataStream unificata.

  • Le applicazioni che utilizzano l' DataSet API devono essere migrate all' DataStream API prima dell'aggiornamento. Consulta la guida alla migrazione di Apache Flink per la conversione DataSet DataStream

Java 11 e Python 3.8 rimossi

  • Il supporto per Java 11 è stato completamente rimosso; Java 17 è il runtime predefinito e consigliato.

  • Il supporto per Python 3.8 è stato rimosso; Python 3.12 è ora l'impostazione predefinita.

Classi di connettori precedenti rimosse

  • Le versioni precedenti SourceFunction e le SinkFunction interfacce sono state sostituite dalle nuove API unificate Source (FLIP-27) e Sink (FLIP-143), che forniscono un supporto migliore per la bounded/unbounded dualità, un migliore coordinamento dei checkpoint e un modello di programmazione più pulito.

  • Per Kinesis Data Streams, usa e da. KinesisStreamsSource KinesisStreamsSink flink-connector-aws-kinesis-streams:6.0.0-2.0

API Scala rimossa

  • L'API Flink Scala è stata rimossa. L'API Java di Flink è ora l'unica API supportata per le JVM-based applicazioni.

  • Se la tua applicazione è scritta in Scala, puoi comunque utilizzare l'API Java di Flink tratta dal codice Scala: la modifica principale è che i Scala-specific wrapper e le conversioni implicite non sono più disponibili. Vedi Aggiornamento delle applicazioni e delle versioni di Flink per dettagli sull'aggiornamento delle applicazioni Scala.

Considerazioni sulla compatibilità degli stati

  • Il serializzatore Kryo aggiornato dalla versione 2.24 alla 5.6 può causare problemi di compatibilità tra stati.

  • I POJO con raccolte (HashMap,,) possono avere problemi di compatibilità tra stati. ArrayList HashSet

  • La serializzazione di Avro e Protobuf non è influenzata.

  • Consulta Guida alla compatibilità dello stato per gli aggiornamenti di Flink 2.2 la sezione per una valutazione dettagliata per valutare il livello di rischio della tua applicazione.

Supporto per runtime e linguaggio

Funzionalità Description Documentazione
Java 17 Runtime Java 17 è ora il runtime predefinito e consigliato; il supporto per Java 11 è stato rimosso. Compatibilità con Java
Supporto per Python 3.12 Python 3.12 ora è supportato; il supporto per Python 3.8 è stato rimosso. PyFlink Documentazione

Gestione dello stato e prestazioni

Funzionalità Description Documentazione
RocksDB 8.10.0 I/O Prestazioni migliorate con l'aggiornamento RocksDB. Backend di stato
Miglioramenti alla serializzazione Serializzatori dedicati per Map, List, Set; Kryo è stato aggiornato dalla versione 2.24 alla 5.6. Serializzazione del tipo

Funzionalità SQL e Table API

Funzionalità Description Documentazione
Tipo di dati VARIANT Supporto nativo per dati semi-strutturati (JSON) senza analisi ripetuta delle stringhe. Tipi di dati
Delta Join Riduce i requisiti statali per lo streaming joins mantenendo solo la versione più recente di ogni chiave; richiede un'infrastruttura gestita dal cliente (ad esempio, Apache Fluss). Si iscrive
StreamingMultiJoinOperator Esegue giunzioni a più vie come un unico operatore, eliminando la materializzazione intermedia. FLIP-516
ProcessTableFunction (PTF) Abilita una logica basata sullo stato e basata sugli eventi direttamente in SQL con stato e timer per chiave. User-Defined Funzioni
Funzione ML_PREDICT Chiama i modelli ML registrati sulle streaming/batch tabelle direttamente da SQL. Richiede al cliente di raggruppare un' ModelProvider implementazione (ad esempio,flink-model-openai). ModelProvider le librerie non vengono fornite da Managed Service for Apache Flink. ML Predict
Modello DDL Definisci i modelli ML come oggetti di catalogo di prima classe utilizzando le istruzioni CREATE MODEL. Dichiarazioni CREATE
Ricerca vettoriale L'API SQL Flink supporta la ricerca nei database vettoriali. Attualmente non è disponibile alcuna VectorSearchTableSource implementazione open source; i clienti devono fornire la propria implementazione. Flink SQL

DataStream Funzionalità API

Funzionalità Description Documentazione
FLIP-27 API di origine Nuova interfaccia sorgente unificata che sostituisce la precedente SourceFunction. Origini
FLIP-143 API Sink Nuova interfaccia sink unificata che sostituisce la precedente SinkFunction. Lavelli
Python asincrono DataStream Non-blocking I/O operazioni nell'API Python utilizzando DataStream . AsyncFunction Asincrono I/O

Quando si esegue l'aggiornamento a Flink 2.2, è inoltre necessario aggiornare le dipendenze del connettore alle versioni compatibili con il runtime di Flink 2.2. I connettori Flink vengono rilasciati indipendentemente dal runtime Flink e non tutti i connettori hanno ancora una versione compatibile con Flink 2.2. La tabella seguente riassume la disponibilità dei connettori di uso comune in Amazon Managed Service for Apache Flink:

Disponibilità dei connettori per Flink 2.2
Connector Versione Flink 1.20 Versione Flink 2.0+ Note
Apache Kafka flink-connector-kafka 3.4.0-1.20 flink-connector-kafka 4.0.0-2.0 Consigliato per Flink 2.2
Kinesis Data Streams (fonte) flink-connector-kinesis 5.0.0-1.20 flink-connector-aws-kinesis-streams 6.0.0-2.0 Consigliato per Flink 2.2
Kinesis Data Streams (sink) flink-connector-aws-kinesis-streams 5.1.0-1.20 connettore-flink-aws-kinesis-streams 6.0.0-2.0 Consigliato per Flink 2.2
Amazon Data Firehose flink-connector-aws-kinesis-firehose 5.1.0-1.20 flink-connector-aws-kinesis-firehose 6.0.0-2.0 Compatibile con Flink 2.0
Amazon DynamoDB flink-connector-dynamodb 5.1.0-1.20 connettore-flink-dynamodb 6.0.0-2.0 Compatibile con Flink 2.0
Amazon SQS flink-connector-sqs 5.1.0-1.20 flink-connector-sqs 6.0.0-2.0 Compatibile con Flink 2.0
FileSystem (S3, HDFS) In dotazione con Flink In dotazione con Flink Integrato nella distribuzione Flink, sempre disponibile
JDBC flink-connector-jdbc 3.3.0-1.20 Non ancora rilasciato per 2.x Nessuna versione compatibile con Flink 2.x disponibile
OpenSearch flink-connector-opensearch 1.2.0-1.19 Non ancora rilasciato per 2.x Nessuna versione compatibile con Flink 2.x disponibile
Elasticsearch Solo connettore precedente Non ancora rilasciato per 2.x Valuta la possibilità di migrare al connettore OpenSearch
Amazon Managed Service per Prometheus flink-connector-prometheus 1.0.0-1.20 Non ancora rilasciato per 2.x Nessuna versione compatibile con Flink 2.x disponibile
  • Se l'applicazione dipende da un connettore che non dispone ancora di una versione Flink 2.x, avete due opzioni: attendere che il connettore rilasci una versione compatibile o valutare se è possibile sostituirlo con un'alternativa (ad esempio, utilizzando il catalogo JDBC o un sink personalizzato).

  • Quando aggiorni le versioni dei connettori, presta attenzione alle modifiche al nome degli artefatti: alcuni connettori sono stati rinominati tra le versioni principali (ad esempio, il connettore Firehose è passato da flink-connector-aws-kinesis-firehose a flink-connector-aws-firehose in alcune versioni intermedie).

  • Consulta sempre la documentazione del connettore Amazon Managed Service for Apache Flink per i nomi esatti degli artefatti e le versioni supportate nel runtime di destinazione.

Le seguenti funzionalità non sono supportate in Amazon Managed Service for Apache Flink 2.2:

  • Tabelle materializzate: istantanee di tabelle interrogabili e gestite in modo continuo.

  • Modifiche alla telemetria personalizzate: report metrici personalizzati e configurazioni di telemetria.

  • ForSt State Backend: archiviazione a stato disaggregata (sperimentale in open source).

  • Java 21: supporto sperimentale in open source, non supportato in Managed Service for Apache Flink.

Amazon Managed Service per Apache Flink Studio

Flink 2.2 in Amazon Managed Service per Apache Flink non supporta le applicazioni Studio. Per ulteriori informazioni, consulta Creazione di un notebook Studio.

Connettore Kinesis EFO

  • Le applicazioni che utilizzano il percorso KinesisStreamsSource with EFO (Enhanced Fan-Out / SubscribeToShard) introdotto nei connettori v5.0.0 e v6.0.0 possono fallire quando gli stream Kinesis vengono sottoposti a resharding. Si tratta di un problema noto nella community. Per ulteriori informazioni, consulta FLINK-37648.

  • Le applicazioni che utilizzano il percorso KinesisStreamsSource with EFO (Enhanced Fan-Out / SubscribeToShard) introdotto nei connettori v5.0.0 e v6.0.0 insieme KinesisStreamsSink possono subire dei deadlock se l'applicazione Flink è sottoposta a contropressione, con conseguente interruzione completa dell'elaborazione dei dati in uno o più casi. TaskManagers Per ripristinare l'applicazione sono necessarie un'operazione di arresto forzato e un'operazione di avvio dell'applicazione. Questo è un caso secondario del problema noto nella community. Per ulteriori informazioni, consulta FLINK-34071.

Amazon Managed Service for Apache Flink supporta aggiornamenti di versione sul posto che preservano la configurazione dell'applicazione, i log, le metriche, i tag e, se lo stato e i file binari sono compatibili, lo stato dell'applicazione. Per istruzioni dettagliate, consulta Aggiornamento a Flink 2.2: guida completa.

Per indicazioni sulla valutazione del rischio di compatibilità degli stati e sulla gestione dello stato incompatibile durante gli aggiornamenti, consulta. Guida alla compatibilità dello stato per gli aggiornamenti di Flink 2.2

Per domande o problemi, consulta Risoluzione dei problemi relativi al servizio gestito per Apache Flink o contatta l' AWS assistenza.