Cette rubrique décrit les principales fonctionnalités et les corrections de bugs incluses dans la version du 19 septembre 2022 de Realtime Compute for Apache Flink.
Présentation
Le 19 septembre 2022, Realtime Compute for Apache Flink a publié une nouvelle version incluant des mises à jour de la plateforme et du moteur, des améliorations des connecteurs, des optimisations des performances ainsi que des corrections de bugs. Cette version comprend VVR-4.0.15 basé sur Apache Flink 1.13 et VVR-6.0.2 basé sur Apache Flink 1.15. Les principales mises à jour sont les suivantes :
Cette version introduit VVR 6.0.2, le premier moteur Flink de niveau entreprise basé sur Apache Flink 1.15. Elle intègre les principales fonctionnalités et optimisations de performances de la communauté open source, notamment des améliorations apportées aux fonctions de table à valeur de fenêtre, aux fonctions CAST, au système de types et aux fonctions JSON.
-
La gestion de l'état constitue une priorité majeure pour nos utilisateurs. Cette version unifie la gestion des checkpoints et des savepoints en une seule fonctionnalité : la gestion des ensembles d'états. Cela améliore considérablement la vitesse de création et de récupération des savepoints, réduit leur taille et renforce le taux de réussite global ainsi que la stabilité.
Par ailleurs, les savepoints ne sont plus supprimés lors de l'annulation d'un déploiement, ce qui modifie le comportement précédent. Désormais, les checkpoints et les savepoints sont distincts et vous pouvez créer et gérer explicitement les savepoints. Outre une meilleure utilisabilité, les optimisations apportées au backend d'état Gemini permettent de réaliser d'importantes économies de coûts. La nouvelle gestion des ensembles d'états peut réduire vos coûts annuels de stockage OSS de 15 % à 40 %. La plateforme vous permet également de démarrer un déploiement à partir d'un savepoint créé par un autre déploiement, ce qui simplifie les tests A/B et autres scénarios de double exécution.
Pour améliorer l'utilisation des ressources, nous avons introduit le réglage planifié. Si votre charge de travail présente des pics et des creux de trafic prévisibles, créez une politique permettant d'ajuster automatiquement les ressources du déploiement à une taille prédéfinie à un moment précis. Cette approche vous aide à gérer la mise à l'échelle des ressources sans intervention manuelle et réduit les coûts opérationnels.
Pour faciliter le diagnostic des déploiements, nous introduisons le score de santé. Cette fonctionnalité analyse les déploiements tout au long de leur cycle de vie, du démarrage à l'exécution, et fournit des informations de diagnostic ainsi que des recommandations pour vous aider à maintenir vos déploiements en streaming.
Pour une meilleure intégration, la plateforme propose désormais un nouvel ensemble d'opérations OpenAPI, vous permettant d'intégrer ses capacités dans vos propres services.
-
Le contrôle des risques en temps réel constitue un cas d'utilisation principal de Flink. Nous avions précédemment proposé en aperçu de nouvelles fonctionnalités de traitement d'événements complexes (CEP) sur des séquences d'événements continues à certains clients sélectionnés, et celles-ci ont été validées avec succès en production.
Dans cette version, nous rendons généralement disponibles une série d'améliorations CEP. Tout d'abord, la capacité très attendue de mise à jour à chaud des règles CEP est désormais disponible. Cela vous permet de mettre à jour les règles pendant les heures de pointe de l'activité sans redémarrer votre déploiement, éliminant ainsi les interruptions de service de dix minutes auxquelles les systèmes de contrôle des risques étaient auparavant confrontés lors des mises à jour de règles, et améliorant considérablement la disponibilité métier. Deuxièmement, nous avons enrichi la syntaxe SQL CEP avec de nouvelles extensions pour améliorer son expressivité. Cela vous permet de convertir des déploiements complexes basés sur l'API DataStream en déploiements SQL plus simples, améliorant l'efficacité du développement et facilitant l'intégration avec les systèmes de traçabilité des données. Enfin, cette version introduit plusieurs nouvelles métriques CEP pour fournir des informations détaillées sur la correspondance des règles.
D'autres optimisations incluent des améliorations de performances. Nous activons désormais automatiquement la séparation clé-valeur pour les opérateurs de jointure de flux doubles dans Flink SQL. Cette optimisation améliore considérablement les performances des déploiements de jointure de flux doubles sans nécessiter de configuration utilisateur. Nous avons également élargi la gamme des versions Hive prises en charge pour Hive Catalog afin d'inclure les versions 2.1.0-2.3.9 et 3.1.0-3.1.3. Pour les connecteurs, nous avons ajouté la prise en charge de la lecture depuis Tablestore et de l'utilisation du connecteur JDBC avec des tables source, des tables de dimension et des tables de résultat.
Nouvelles fonctionnalités
Fonctionnalité | Description | Documentation |
Gestion des ensembles d'états | La gestion des ensembles d'états découple la gestion de l'état des opérations de démarrage et d'arrêt pour tous les déploiements Flink avec état. Les savepoints ne sont plus supprimés lorsqu'un déploiement est arrêté. Vous pouvez utiliser une page de gestion dédiée pour créer et supprimer des savepoints selon une planification. | |
Réglage planifié | Pour les déploiements Flink présentant des pics et des creux de trafic prévisibles, définissez des politiques de planification personnalisées. Aux moments spécifiés, les ressources du déploiement sont automatiquement ajustées à une taille prédéfinie pour gérer les fluctuations de trafic, éliminant ainsi le besoin d'une mise à l'échelle manuelle. | |
Score de santé | La fonctionnalité de score de santé applique des règles expertes pour détecter les problèmes lors du démarrage et de l'exécution du déploiement, en fournissant des recommandations actionnables. Cette fonctionnalité vous aide à mieux comprendre l'état de vos déploiements et à ajuster les paramètres en conséquence. | |
Amélioration de l'autorisation des membres | Le processus d'autorisation est amélioré : au lieu de saisir manuellement les informations utilisateur, sélectionnez parmi une liste de tous les utilisateurs RAM lors de l'octroi des autorisations. | |
Traitement dynamique des événements complexes (CEP) | Le CEP offre des capacités de correspondance de motifs pour les flux de données en temps réel. Cette version s'appuie sur Flink CEP open source en vous permettant d'externaliser les règles de déploiement dans une base de données afin qu'elles puissent être chargées dynamiquement. Cette fonctionnalité est exposée via l'API DataStream. | |
Amélioration du SQL CEP | L'instruction MATCH_RECOGNIZE vous permet de décrire les règles CEP à l'aide de SQL. Cette version améliore l'instruction MATCH_RECOGNIZE de Flink open source avec de nouvelles capacités, telles que la sortie des correspondances expirées et la prise en charge de De plus, de nouvelles métriques ont été introduites :
| |
Prise en charge de la synchronisation de base de données vers Kafka | Lorsque vous utilisez cette fonctionnalité, les données sont synchronisées vers une table Upsert Kafka correspondante. Vous pouvez utiliser directement la table dans Kafka au lieu de la table MySQL, ce qui réduit la charge sur le service MySQL provenant de plusieurs déploiements. | |
Définir des tables partitionnées dans les tables de résultat Hologres avec DDL | Vous pouvez utiliser PARTITION BY pour définir une table partitionnée lors de la création d'une table de résultat Hologres. | |
Définir le délai d'expiration pour les requêtes asynchrones dans les tables de dimension Hologres | En définissant le paramètre | |
Définir les propriétés de table lors de la création de tables avec Hologres Catalog | La définition de propriétés de table appropriées peut aider le système à organiser et interroger les données efficacement. Lorsque vous utilisez Hologres Catalog pour créer une table, vous pouvez désormais définir les propriétés de table physique dans la clause WITH. | |
Le connecteur sink MaxCompute prend en charge le type Binary |
| |
Hive Catalog prend en charge davantage de versions Hive | Cette version prend en charge Hive 2.1.0-2.3.9 et 3.1.0-3.1.3. | |
Connecteur source Tablestore publié | Prend en charge la lecture des journaux incrémentiels depuis Tablestore. | |
Connecteur JDBC publié | Le connecteur JDBC communautaire est désormais intégré. | |
Le parallélisme d'une table source Message Queue for Apache RocketMQ peut dépasser le nombre de partitions de topic | Ce mode vous permet de pré-allouer des ressources pour d'éventuelles augmentations des partitions de topic avant le début de la consommation. | |
Définir la clé de message pour les tables de résultat Message Queue for Apache RocketMQ | Vous pouvez désormais définir la clé de message lors de l'écriture dans Message Queue for Apache RocketMQ. | |
Prise en charge du catalogue AnalyticDB for MySQL | Avec ce catalogue, vous pouvez lire directement les métadonnées depuis AnalyticDB for MySQL sans enregistrer manuellement les tables AnalyticDB for MySQL. Cela améliore l'efficacité du développement et garantit l'exactitude des données. |
Optimisations des performances
-
Cette version introduit un format de savepoint natif, qui résout les problèmes de délai d'expiration qui se produisaient auparavant avec les savepoints au format standard pour les déploiements comportant de grands états. Cela améliore considérablement la stabilité globale du déploiement.
Métrique
Amélioration
Temps de finalisation du savepoint
Une amélioration moyenne de 5 à 10 fois, le ratio augmentant à mesure que la taille de l'état incrémentiel diminue. Dans certains déploiements typiques, l'amélioration peut atteindre 100 fois.
Temps de récupération du déploiement
Une amélioration moyenne d'environ 5 fois, le ratio augmentant à mesure que la taille de l'état augmente.
Frais généraux d'espace du savepoint
Une réduction moyenne des frais généraux d'espace de 2 fois, le ratio augmentant à mesure que la taille de l'état augmente.
Frais généraux réseau du savepoint
Une réduction moyenne des frais généraux réseau de 5 à 10 fois, le ratio augmentant à mesure que la taille de l'état incrémentiel diminue.
L'opérateur de jointure de flux doubles infère désormais automatiquement quand activer la séparation clé-valeur pour optimiser les performances. Pour les déploiements SQL, l'opérateur de jointure de flux doubles analyse automatiquement les caractéristiques du déploiement et active la séparation clé-valeur pour optimiser les performances. Lors des tests de performance pour des scénarios typiques, les performances moyennes se sont améliorées de plus de 40 %. Pour plus d'informations, consultez Optimiser les performances élevées de Flink SQL et Configurer les backends d'état de niveau entreprise.
Le démarrage du déploiement est désormais 15 % plus rapide en moyenne.
Corrections de bugs
Correction d'un problème où l'heure de modification d'un déploiement était mise à jour incorrectement.
Correction d'un problème où l'état de certains déploiements ne pouvait pas être déterminé après avoir été suspendu et redémarré.
Correction d'un problème empêchant le téléchargement local de fichiers JAR depuis Alibaba Finance Cloud.
Correction d'un problème où les ressources totales utilisées par un déploiement en cours d'exécution étaient incohérentes avec les statistiques affichées sur la page.
Correction d'un problème d'échec de navigation dans les journaux de diagnostic du déploiement.
Correction d'une erreur lors de la lecture directe d'une table upsert Kafka depuis un catalogue Kafka.
Correction d'une NullPointerException lors de l'utilisation de résultats intermédiaires dans des opérations imbriquées avec plusieurs fonctions définies par l'utilisateur (UDF).
Correction de problèmes dans mysql-cdc, notamment un fractionnement anormal des chunks, des erreurs de mémoire insuffisante (OOM) et des fuseaux horaires incohérents entre les données initiales et incrémentielles. Pour plus d'informations, consultez Table source MySQL CDC.