View a markdown version of this page

Amazon Managed Service para Apache Flink 2.2 - Managed Service para Apache Flink

Las traducciones son generadas a través de traducción automática. En caso de conflicto entre la traducción y la version original de inglés, prevalecerá la version en inglés.

Amazon Managed Service para Apache Flink 2.2

Amazon Managed Service para Apache Flink ahora es compatible con la versión 2.2 de Apache Flink. Esta es la primera actualización importante de la versión del servicio. En esta página se describen las funciones introducidas en Flink 2.2, junto con consideraciones importantes para actualizar desde Flink 1.x.

nota

Flink 2.2 presenta cambios importantes que requieren una planificación cuidadosa. Consulta la lista completa de cambios importantes y obsoletos que aparece a continuación y Guía de compatibilidad estatal para las actualizaciones de Flink 2.2 antes de actualizar desde la versión 1.x.

Amazon Managed Service para Apache Flink 2.2 introduce cambios de comportamiento que pueden dañar las aplicaciones existentes tras la actualización. Revíselos detenidamente junto con los cambios en la API de Flink en la siguiente sección.

Gestión programática de la configuración

Eliminación de métricas

  • La fullRestarts métrica se ha eliminado en Flink 2.2. En su lugar, usa la numRestarts métrica.

  • La bytesRequestedPerFetch métrica del conector KDS se eliminó en la versión 6.0.0 AWS del conector Flink (la única versión del conector compatible con Flink 2.2).

  • uptimeTanto la métrica como la downtime métrica están marcadas como obsoletas en Flink 2.2 y se eliminarán pronto. uptimeSustitúyala por la nueva métrica. runningTime downtimeSustitúyalo por uno o más de restartingTimecancellingTime, yfailingTime.

  • Consulta la página de métricas y dimensiones para ver la lista completa de las métricas compatibles.

Non-Credential Llamadas de IMDS bloqueadas

  • Los AWS SDK () y DefaultCredentialsProvider (/latest/meta-data/iam/security-credentials/) utilizan estos puntos de conexión permitidos para configurar automáticamente las credenciales y DefaultAwsRegionProviderChain la región de la aplicación. /latest/dynamic/instance-identity/document

  • Las aplicaciones que utilizan las funciones del SDK de AWS que se basan en llamadas IMDS sin credenciales (comoEC2MetadataUtils.getInstanceId(), EC2MetadataUtils.getInstanceType()EC2MetadataUtils.getLocalHostName(), oEC2MetadataUtils.getAvailabilityZone()) recibirán errores HTTP 4xx al intentar realizar estas llamadas.

  • Si su aplicación usa IMDS, por ejemplo, metadatos u otra información fuera de las rutas permitidas, refactorice el código para usar variables de entorno o la configuración de la aplicación en su lugar.

Read-Only Sistema de archivos raíz

  • Para mejorar la seguridad, cualquier dependencia fuera de la /tmp cual se encuentre el directorio de trabajo predeterminado de Flink dará como resultado:. java.io.FileNotFoundException: /{path}/{filename} (Read-only file system)

  • Las dependencias del sistema de archivos pueden originarse directamente en su código o indirectamente en las bibliotecas incluidas en sus dependencias. Anula las dependencias directas del sistema de archivos para incluirlas en tu código. /tmp/ Para las dependencias indirectas del sistema de archivos de las bibliotecas, usa las anulaciones de configuración de las bibliotecas para redirigir las operaciones del sistema de archivos a. /tmp/

A continuación se muestra un resumen de los cambios más importantes y las obsolescencias introducidos en Managed Service para Apache Flink 2.2. Consulte las notas de la versión 2.0 de Apache Flink para ver las notas completas de la versión de Apache Flink 2.0, en las que se presentan estos cambios importantes.

Eliminación de la API y el lenguaje de Flink

DataSet API eliminada

Se eliminaron Java 11 y Python 3.8

  • Se eliminó por completo la compatibilidad con Java 11; Java 17 es el tiempo de ejecución predeterminado y recomendado.

  • Se ha eliminado la compatibilidad con Python 3.8; Python 3.12 es ahora el predeterminado.

Se han eliminado las clases de conectores antiguas

  • SinkFunctionLas interfaces SourceFunction y las antiguas han sido reemplazadas por las nuevas API unificadas Source (FLIP-27) y Sink (FLIP-143), que ofrecen un mejor soporte para la bounded/unbounded dualidad, una mejor coordinación de los puntos de control y un modelo de programación más limpio.

  • Para Kinesis Data Streams, utilice y desde. KinesisStreamsSource KinesisStreamsSink flink-connector-aws-kinesis-streams:6.0.0-2.0

Se ha eliminado la API de Scala

  • Se ha eliminado la API de Flink Scala. La API Java de Flink es ahora la única API compatible para las aplicaciones. JVM-based

  • Si tu aplicación está escrita en Scala, puedes seguir usando la API Java de Flink a partir del código de Scala. La principal diferencia es que los contenedores y las conversiones implícitas Scala-specific ya no están disponibles. Consulte Actualización de aplicaciones y versiones de Flink para obtener más información sobre la actualización de sus aplicaciones de Scala.

Consideraciones sobre la compatibilidad entre estados

  • La actualización del serializador Kryo de la versión 2.24 a la 5.6 puede provocar problemas de compatibilidad de estados.

  • Los POJO con colecciones (HashMap,ArrayList,HashSet) pueden tener problemas de compatibilidad entre estados.

  • La serialización de Avro y Protobuf no se ve afectada.

  • Consulte Guía de compatibilidad estatal para las actualizaciones de Flink 2.2 para obtener una evaluación detallada para determinar el nivel de riesgo de su aplicación.

Soporte de lenguaje y tiempo de ejecución

Característica Description (Descripción) Documentación
Tiempo de ejecución de Java 17 Java 17 es ahora el tiempo de ejecución predeterminado y recomendado; se ha eliminado la compatibilidad con Java 11. Compatibilidad con Java
Soporte para Python 3.12 Ahora es compatible con Python 3.12; se ha eliminado el soporte para Python 3.8. PyFlink Documentación

Gestión del estado y rendimiento

Característica Description (Descripción) Documentación
RockSDB 8.10.0 I/O Rendimiento mejorado con la actualización de RockSDB. Backends estatales
Mejoras en la serialización Serializadores dedicados para Map, List y Set; Kryo se actualizó de la 2.24 a la 5.6. Escriba: Serialización

Características de la API de SQL y Table

Característica Description (Descripción) Documentación
Tipo de datos VARIANT Soporte nativo para datos semiestructurados (JSON) sin análisis repetido de cadenas. Tipos de datos
Delta Join Reduce los requisitos estatales para las uniones por streaming al mantener solo la versión más reciente de cada clave; requiere una infraestructura gestionada por el cliente (por ejemplo, Apache Fluss). Se une
StreamingMultiJoinOperator Ejecuta uniones multidireccionales como un solo operador, lo que elimina la materialización intermedia. FLIP-516
ProcessTableFunction (PTF) Permite la lógica con estado y basada en eventos directamente en SQL con temporizadores y estados por clave. User-Defined Funciones
Función ML_PREDICT Llame a los modelos de ML registrados en streaming/batch las tablas directamente desde SQL. Requiere que el cliente agrupe una ModelProvider implementación (p. ej.,flink-model-openai). ModelProvider Managed Service no distribuye las bibliotecas para Apache Flink. ML Predict
Modelo DDL Defina los modelos de ML como objetos de catálogo de primera clase mediante las instrucciones CREATE MODEL. Declaraciones CREATE
Búsqueda vectorial La API SQL de Flink admite la búsqueda en bases de datos vectoriales. Actualmente no hay ninguna VectorSearchTableSource implementación de código abierto disponible; los clientes deben proporcionar su propia implementación. SQL de Flink

DataStream Funciones de la API

Característica Description (Descripción) Documentación
FLIP-27 API de origen Nueva interfaz de código fuente unificada que reemplaza a la antigua SourceFunction. Orígenes
FLIP-143 API Sink Nueva interfaz de sumidero unificada que reemplaza a la antigua SinkFunction. Lavabos
Python asíncrono DataStream Non-blocking I/O operaciones en la DataStream API de Python utilizando. AsyncFunction Asincrónico I/O

Al actualizar a Flink 2.2, también debes actualizar las dependencias de tus conectores a versiones que sean compatibles con el tiempo de ejecución de Flink 2.2. Los conectores Flink se lanzan independientemente del tiempo de ejecución de Flink, y no todos los conectores tienen todavía una versión compatible con Flink 2.2. En la siguiente tabla se resume la disponibilidad de los conectores más utilizados en Amazon Managed Service para Apache Flink:

Disponibilidad de conectores para Flink 2.2
Connector Versión 1.20 de Flink Versión Flink 2.0+ Notas
Apache Kafka flink-connector-kafka 3.4.0-1.20 flink-connector-kafka 4.0.0-2.0 Recomendado para Flink 2.2
Kinesis Data Streams (fuente) flink-connector-kinesis 5.0.0-1.20 flink-connector-aws-kinesis-streams 6.0.0-2.0 Recomendado para Flink 2.2
Kinesis Data Streams (sumidero) flink-connector-aws-kinesis-streams 5.1.0-1.20 flink-connector-aws-kinesis-streams 6.0.0-2.0 Recomendado para 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 Compatible con Flink 2.0
Amazon DynamoDB conector flink-dynamodb 5.1.0-1.20 conector flink-dynamodb 6.0.0-2.0 Compatible con Flink 2.0
Amazon SQS flink-connector-sqs 5.1.0-1.20 flink-connector-sqs 6.0.0-2.0 Compatible con Flink 2.0
FileSystem (S3, HDFS) Incluido con Flink Incluido con Flink Integrado en la distribución de Flink, siempre disponible
JDBC flink-connector-jdbc 3.3.0-1.20 Aún no se ha publicado para 2.x No hay ninguna versión compatible con Flink 2.x
OpenSearch flink-connector-opensearch 1.2.0-1.19 Aún no se ha publicado para 2.x No hay ninguna versión compatible con Flink 2.x
Elasticsearch Solo conector Legacy Aún no se ha publicado para la versión 2.x Considere la posibilidad de migrar al conector OpenSearch
Servicio administrado por Amazon para Prometheus flink-connector-prometheus 1.0.0-1.20 Aún no se ha publicado para 2.x No hay ninguna versión compatible con Flink 2.x
  • Si tu aplicación depende de un conector que aún no tenga una versión 2.x de Flink, tienes dos opciones: esperar a que el conector publique una versión compatible o evaluar si puedes reemplazarlo por una alternativa (por ejemplo, usar el catálogo de JDBC o un receptor personalizado).

  • Al actualizar las versiones de los conectores, preste atención a los cambios en el nombre de los artefactos: algunos conectores cambiaron de nombre entre las versiones principales (por ejemplo, el conector Firehose cambió de flink-connector-aws-kinesis-firehose a flink-connector-aws-firehose en algunas versiones intermedias).

  • Consulta siempre la documentación del conector Apache Flink de Amazon Managed Service para ver los nombres exactos de los artefactos y las versiones compatibles en el tiempo de ejecución de destino.

Amazon Managed Service para Apache Flink 2.2 no admite las siguientes funciones:

  • Tablas materializadas: instantáneas de tablas consultables y mantenidas de forma continua.

  • Cambios de telemetría personalizados: indicadores métricos y configuraciones de telemetría personalizados.

  • ForSt Backend estatal: almacenamiento estatal desagregado (experimental en código abierto).

  • Java 21: soporte experimental en código abierto, no compatible con Managed Service para Apache Flink.

Servicio gestionado de Amazon para Apache Flink Studio

Flink 2.2 del Amazon Managed Service para Apache Flink no es compatible con las aplicaciones de Studio. Para obtener más información, consulte Creación de un cuaderno de Studio.

Conector Kinesis EFO

  • Las aplicaciones que utilizan la ruta KinesisStreamsSource con EFO (Enhanced Fan-Out / SubscribeToShard) introducida en los conectores v5.0.0 y v6.0.0 pueden fallar cuando las transmisiones de Kinesis se vuelven a fragmentar. Se trata de un problema conocido en la comunidad. Para obtener más información, consulte FLINK-37648.

  • Las aplicaciones que utilizan la KinesisStreamsSource ruta EFO (Enhanced Fan-Out / SubscribeToShard) introducida en los conectores v5.0.0 y v6.0.0, junto con la ruta EFO (Enhanced/), KinesisStreamsSink pueden bloquearse si la aplicación Flink sufre una contrapresión, lo que provoca que se detenga por completo el procesamiento de datos en uno o más. TaskManagers Para recuperar la aplicación es necesaria una operación de parada forzada y una operación de inicio de la aplicación. Este es un subcaso del problema conocido en la comunidad. Para obtener más información, consulte FLINK-34071.

Amazon Managed Service para Apache Flink admite actualizaciones de versiones in situ que preservan la configuración, los registros, las métricas, las etiquetas y, si el estado y los binarios son compatibles, el estado de la aplicación. Para obtener instrucciones paso a paso, consulte Actualización a Flink 2.2: guía completa.

Para obtener orientación sobre cómo evaluar el riesgo de compatibilidad entre estados y gestionar los estados incompatibles durante las actualizaciones, consulte. Guía de compatibilidad estatal para las actualizaciones de Flink 2.2

Si tiene preguntas o problemas, consulte el Resolución de problemas de Managed Service para Apache Flink o póngase en contacto con el servicio de AWS asistencia.