Synchronisez des données vers un catalogue DLF via Iceberg REST en utilisant Flink CDC sur Realtime Compute for Apache Flink.
Prérequis
Vous disposez d'un espace de travail entièrement géré pour Realtime Compute for Apache Flink. Si vous n'en avez pas encore créé, consultez la rubrique Activer Realtime Compute for Apache Flink.
Vérifiez que votre espace de travail Realtime Compute for Apache Flink et DLF se trouvent dans la même région. Vous devez également ajouter le VPC de l'espace de travail à la liste d'autorisation de DLF. Pour plus d'informations, consultez la section Configurer une liste d'autorisation VPC.
Limites
La connectivité Iceberg REST à DLF nécessite le moteur Realtime Compute for Apache Flink VVR 11.6.0 ou une version ultérieure.
Enregistrer le catalogue DLF dans Flink
L'enregistrement d'un catalogue dans Flink crée un mappage vers votre catalogue DLF. La création ou la suppression du catalogue dans Flink n'affecte pas les données réelles stockées dans DLF. Toutes les tables créées dans le catalogue DLF via Iceberg REST sont des tables Iceberg.
Connectez-vous à la console de gestion Realtime Compute for Apache Flink.
Dans la colonne Actions de votre espace de travail, cliquez sur Console.
Dans le volet de navigation de gauche, cliquez sur Development > Scripts.
-
Créez un nouveau script, puis collez l'instruction SQL suivante dans l'éditeur SQL.
CREATE CATALOG `catalog_name` WITH ( 'type' = 'iceberg', 'catalog-type' = 'rest', 'uri' = 'http://{region-id}-vpc.dlf.aliyuncs.com/iceberg', 'warehouse' = 'iceberg_test', 'rest.signing-region' = '{region-id}', 'io-impl' = 'org.apache.iceberg.rest.DlfFileIO' );Remplacez
{region-id}par l'ID de région de votre catalogue DLF, par exemplecn-hangzhououap-southeast-1. Pour connaître toutes les régions prises en charge et leurs valeurs d'endpoint, consultez la rubrique Régions et endpoints. Dans le coin inférieur droit, cliquez sur Environment, sélectionnez un cluster de session exécutant VVR 11.2.0 ou une version ultérieure, puis exécutez l'instruction SQL.
Le tableau suivant décrit les options de configuration.
| Option | Description | Required | Example |
|---|---|---|---|
type |
Type de catalogue. Définissez la valeur sur iceberg. |
Oui | iceberg |
catalog-type |
Type de catalogue. Définissez la valeur sur rest. |
Oui | rest |
token.provider |
Fournisseur de jetons pour l'authentification DLF. Définissez la valeur sur dlf. |
Oui | dlf |
uri |
Endpoint Iceberg REST de votre catalogue DLF. Utilisez le format http://{region-id}-vpc.dlf.aliyuncs.com/iceberg. Pour les valeurs spécifiques à chaque région, consultez la rubrique Régions et endpoints. |
Oui | http://ap-southeast-1-vpc.dlf.aliyuncs.com/iceberg |
warehouse |
Nom de votre catalogue DLF. | Oui | iceberg_test |
rest.signing-region |
ID de région de votre catalogue DLF. Pour obtenir les ID de région, consultez la page Régions et endpoints. | Oui | ap-southeast-1 |
io-impl |
Implémentation FileIO pour DLF. Définissez la valeur sur org.apache.iceberg.rest.DlfFileIO. |
Oui | org.apache.iceberg.rest.DlfFileIO |
Configurer Flink CDC pour utiliser le catalogue
Créez une tâche d'ingestion de données : Développer une tâche d'ingestion de données Flink CDC.
Si vous disposez déjà d'un mappage de catalogue Flink, Réutilisez un catalogue existant pour obtenir les informations de connexion et ajoutez cette configuration Sink :
sink:
type: iceberg
using.built-in-catalog: catalog_name
Exemples de configuration
Modèles courants de synchronisation de données utilisant des tâches YAML Flink CDC :
Synchroniser une base de données MySQL complète vers DLF
Cette tâche synchronise l'intégralité d'une base de données MySQL vers DLF :
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: mysql_test.\.*
server-id: 8601-8604
# (Optional) Sync tables created during incremental phase.
scan.binlog.newly-added-table.enabled: true
# (Optional) Sync table and field comments.
include-comments.enabled: true
# (Optional) Prioritize unbounded chunks to prevent TaskManager OOM errors.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
# (Optional) Speed up reads by parsing only matched tables.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: iceberg
using.built-in-catalog: catalog_name
Paramètres optionnels recommandés pour les sources MySQL (MySQL) :
-
Paramètre : scan.binlog.newly-added-table.enabled
Fonction : Synchronise les tables créées pendant la phase incrémentielle.
-
Paramètre : include-comments.enabled
Fonction : Synchronise les commentaires de table et de champ.
-
Paramètre : scan.incremental.snapshot.unbounded-chunk-first.enabled
Fonction : Empêche les erreurs OOM (Out Of Memory) du TaskManager.
-
Paramètre : scan.only.deserialize.captured.tables.changelog.enabled: true
Fonction : Accélère les lectures en analysant uniquement les tables correspondantes.
Écrire dans une table DLF partitionnée
Spécifiez les clés de partition avec le paramètre partition-keys (Référence de développement des tâches d'ingestion de données Flink CDC) :
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: mysql_test.\.*
server-id: 8601-8604
# (Optional) Sync tables created during incremental phase.
scan.binlog.newly-added-table.enabled: true
# (Optional) Sync table and field comments.
include-comments.enabled: true
# (Optional) Prioritize unbounded chunks to prevent TaskManager OOM errors.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
# (Optional) Speed up reads by parsing only matched tables.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: iceberg
using.built-in-catalog: catalog_name
transform:
- source-table: mysql_test.tbl1
# (Optional) Set partition keys.
partition-keys: id,pt
- source-table: mysql_test.tbl2
partition-keys: id,pt
Écrire dans une table DLF en ajout seul (append-only)
Pour implémenter des suppressions logicielles (convertir les opérations DELETE en INSERT en aval), configurez la tâche comme suit (Référence de développement des tâches d'ingestion de données Flink CDC) :
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: mysql_test.\.*
server-id: 8601-8604
# (Optional) Sync tables created during incremental phase.
scan.binlog.newly-added-table.enabled: true
# (Optional) Sync table and field comments.
include-comments.enabled: true
# (Optional) Prioritize unbounded chunks to prevent TaskManager OOM errors.
scan.incremental.snapshot.unbounded-chunk-first.enabled: true
# (Optional) Speed up reads by parsing only matched tables.
scan.only.deserialize.captured.tables.changelog.enabled: true
sink:
type: iceberg
using.built-in-catalog: catalog_name
transform:
- source-table: mysql_test.tbl1
# (Optional) Set partition keys.
partition-keys: id,pt
# (Optional) Implement soft delete.
projection: \*, __data_event_type__ AS op_type
converter-after-transform: SOFT_DELETE
- source-table: mysql_test.tbl2
# (Optional) Set partition keys.
partition-keys: id,pt
# (Optional) Implement soft delete.
projection: \*, __data_event_type__ AS op_type
converter-after-transform: SOFT_DELETE
L'ajout de
__data_event_type__à la projection écrit le type d'événement de modification sous forme de nouveau champ en aval. La définition deconverter-after-transformsurSOFT_DELETEconvertit les suppressions en insertions, enregistrant ainsi tous les événements de modification. Référence de développement des tâches d'ingestion de données Flink CDC.
Synchroniser des données Kafka CDC en temps réel vers DLF
Cet exemple synchronise les données de modification de deux tables (clients et produits) depuis un topic d'inventaire Kafka au format Debezium JSON vers DLF :
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: ${kafka.bootstrap.servers}
topic: inventory
scan.startup.mode: earliest-offset
value.format: debezium-json
debezium-json.distributed-tables: true
sink:
type: iceberg
using.built-in-catalog: catalog_name
# Debezium JSON lacks primary key info. Add it manually.
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
La source Kafka prend en charge les formats
canal-json,debezium-json(par défaut) etjson.-
Lorsque vous utilisez
debezium-json, ajoutez manuellement une clé primaire à l'aide d'une règle de transformation, car les messages Debezium JSON ne contiennent pas d'informations sur la clé primaire :transform: - source-table: \.*.\.* projection: \* primary-keys: id Si les données d'une seule table s'étendent sur plusieurs partitions, ou si des tables situées dans différentes partitions doivent être fusionnées, définissez
debezium-json.distributed-tablesoucanal-json.distributed-tablessurtrue.La source Kafka prend en charge plusieurs stratégies d'inférence de schéma via le paramètre
schema.inference.strategy. Kafka.
Synchroniser des journaux Kafka en temps réel vers DLF
Pour les données JSON personnalisées dans Kafka, Flink CDC gère automatiquement l'inférence des types de données, l'inférence de schéma et l'évolution de schéma.
Cet exemple synchronise une seule table de journaux JSON depuis le topic d'inventaire vers DLF :
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: ${kafka.bootstrap.servers}
topic: inventory
scan.startup.mode: earliest-offset
value.format: json
# (Optional) Recursively flatten nested columns in JSON data.
json.infer-schema.flatten-nested-columns.enable: true
# (Optional) Skip first 100 parsing errors. Job fails if errors exceed 100.
ingestion.ignore-errors: true
ingestion.error-tolerance.max-count: 100
sink:
type: iceberg
using.built-in-catalog: catalog_name
# Add primary key to table.
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
# Write all inventory topic data to test_database.inventory.
route:
- source-table: inventory
sink-table: test_database.inventory
pipeline:
# (Optional) Log dirty data that causes processing exceptions.
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger
Cet exemple synchronise plusieurs tables de journaux JSON depuis le topic d'inventaire, en utilisant les champs databaseName et tableName pour identifier les tables :
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: ${kafka.bootstrap.servers}
topic: inventory
scan.startup.mode: earliest-offset
value.format: json
# (Optional) Recursively flatten nested columns in JSON data.
json.infer-schema.flatten-nested-columns.enable: true
# Use databaseName and tableName field values as database and table names.
json.decode.parser-table-id.fields: databaseName,tableName
# (Optional) Skip first 100 parsing errors. Job fails if errors exceed 100.
ingestion.ignore-errors: true
ingestion.error-tolerance.max-count: 100
sink:
type: iceberg
using.built-in-catalog: catalog_name
# Add primary key to tables.
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
# Write ods.inventory, ods.customer, and ods.user to test_database.inventory, test_database.customer, and test_database.user respectively.
route:
- source-table: ods.inventory
sink-table: test_database.inventory
- source-table: ods.customer
sink-table: test_database.customer
- source-table: ods.user
sink-table: test_database.user
pipeline:
# (Optional) Log dirty data that causes processing exceptions.
dirty-data.collector:
name: Logger Dirty Data Collector
type: logger
Stratégies d'inférence et d'évolution de schéma JSON : Stratégies d'analyse de schéma et de synchronisation des modifications.
Référence de configuration complète : Référence de développement des tâches d'ingestion de données Flink CDC.