Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Connecteur YAML Hologres pour l'ingestion de données

Dernière mise à jour :Aug 09, 2026

Cette rubrique explique comment utiliser le connecteur Hologres pour synchroniser les données dans un déploiement d'ingestion de données YAML.

Contexte

Hologres est un moteur d'entrepôt de données en temps réel qui prend en charge l'écriture, la mise à jour et l'analyse en temps réel à grande échelle. Hologres est compatible avec le protocole PostgreSQL et prend en charge le langage SQL standard. Il permet le traitement analytique en ligne (OLAP) et des requêtes ad hoc sur des pétaoctets de données, tout en assurant une diffusion des données avec une haute concurrence et une faible latence. Hologres s'intègre à MaxCompute, Realtime Compute for Apache Flink et DataWorks pour offrir une solution complète d'entrepôt de données en ligne et hors ligne. Le tableau suivant décrit les capacités du connecteur YAML Hologres.

Élément

Description

Type de table

Sink

Mode d'exécution

Modes streaming et batch

Format des données

N/A

Métrique

  • numRecordsOut

  • numRecordsOutPerSecond

Remarque

Pour plus de détails, consultez la section Métriques.

Type d'API

YAML

Mises à jour ou suppressions dans une table sink

Pris en charge

Fonctionnalités

Fonctionnalité

Description

Synchronisation complète de la base de données

Synchronise en temps réel les données complètes et incrémentielles depuis une base de données entière ou plusieurs tables vers les tables sink correspondantes.

Synchronisation des modifications de schéma

Synchronise en temps réel les modifications de schéma (telles que l'ajout, la suppression ou le renommage de colonnes) depuis les tables sources vers les tables sink correspondantes.

Synchronisation de bases de données et de tables fragmentées

Utilisez une expression régulière pour faire correspondre les tables sources dans plusieurs bases de données fragmentées par nom. Les données de ces tables sont ensuite fusionnées et synchronisées vers les tables sink aval portant le même nom.

Écriture dans une table partitionnée

Écrit les données d'une table amont dans une table partitionnée Hologres.

Mappage des types de données

Mappe les types de données amont vers des types de données Hologres plus larges à l'aide de plusieurs stratégies.

Syntaxe

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}

Paramètres

Paramètre

Description

Type

Obligatoire

Valeur par défaut

Remarques

type

Le type de sink.

String

Oui

Aucune

La valeur doit être hologres.

name

Le nom du sink.

String

Non

Aucune

N/A.

dbname

Le nom de la base de données.

String

Oui

Aucune

N/A.

username

Le nom d'utilisateur pour l'accès à la base de données. Utilisez l'ID AccessKey de votre compte Alibaba Cloud.

String

Oui

Aucune

Pour plus d'informations, consultez la section Comment afficher l'ID AccessKey et la clé secrète AccessKey ?

Important

Pour éviter d'exposer votre AccessKey, spécifiez ses valeurs à l'aide de variables. Pour plus d'informations, consultez la section Variables de projet.

password

Le mot de passe pour l'accès à la base de données. Utilisez la clé secrète AccessKey de votre compte Alibaba Cloud.

String

Oui

Aucune

endpoint

Le point de terminaison du service Hologres.

String

Oui

Aucune

Pour plus d'informations, consultez la section Points de terminaison d'accès.

jdbcRetryCount

Le nombre de tentatives pour les opérations d'écriture et de requête en cas d'échec de connexion.

Integer

Non

10

N/A.

jdbcRetrySleepInitMs

Le temps d'attente fixe pour chaque tentative.

Long

Non

1000

Unité : millisecondes. Le temps d'attente réel pour une nouvelle tentative est calculé à l'aide de la formule suivante : jdbcRetrySleepInitMs + retry * jdbcRetrySleepStepMs.

jdbcRetrySleepStepMs

Le temps d'attente incrémentiel pour chaque tentative.

Long

Non

5000

Unité : millisecondes. Le temps d'attente réel pour une nouvelle tentative est calculé à l'aide de la formule suivante : jdbcRetrySleepInitMs + retry * jdbcRetrySleepStepMs.

jdbcConnectionMaxIdleMs

Le temps d'inactivité maximal pour une connexion JDBC.

Long

Non

60000

Unité : millisecondes. Si une connexion reste inactive plus longtemps que cette durée, elle est déconnectée et libérée.

jdbcMetaCacheTTL

Le délai d'expiration des informations TableSchema mises en cache localement.

Long

Non

60000

Unité : millisecondes.

jdbcMetaAutoRefreshFactor

Un facteur qui détermine le déclencheur d'actualisation du cache. Le système actualise automatiquement le cache si sa durée de vie restante est inférieure au temps de déclenchement.

Integer

Non

4

La durée de vie restante du cache est calculée comme suit : Durée de vie restante du cache = Délai d'expiration du cache - Temps d'activité du cache. Après une actualisation automatique, le temps d'activité du cache est réinitialisé à 0.

Le temps de déclenchement est calculé à l'aide de la formule : jdbcMetaCacheTTL / jdbcMetaAutoRefreshFactor.

mutatetype

Le mode d'écriture des données.

String

Non

INSERT_OR_UPDATE

Si une clé primaire est définie pour la table physique Hologres, le sink Hologres garantit une sémantique exactly-once basée sur la clé primaire. Si des données avec une clé primaire en double arrivent, le paramètre mutatetype détermine comment la table sink est mise à jour. Le paramètre mutatetype prend en charge les valeurs suivantes :

  • INSERT_OR_IGNORE : Conserve la première occurrence des données et ignore toutes les suivantes.

  • INSERT_OR_REPLACE : Les nouvelles données remplacent toute la ligne existante.

  • INSERT_OR_UPDATE : Met à jour un sous-ensemble de colonnes pour une ligne existante. Par exemple, une table possède quatre colonnes : a, b, c et d, et a est la clé primaire (PK). Si vous écrivez des données uniquement pour les colonnes a et b dans Hologres et qu'une ligne avec la même PK existe déjà, le système met à jour uniquement la colonne b. Les colonnes c et d restent inchangées.

createparttable

Indique s'il faut créer automatiquement les partitions inexistantes lors de l'écriture dans une table partitionnée.

Boolean

Non

false

N/A.

sink.delete-strategy

Spécifie comment traiter les messages de rétractation.

String

Non

Aucune

Valeurs valides :

  • IGNORE_DELETE : Ignore les messages Update Before et Delete. Cette valeur convient aux scénarios où vous devez uniquement insérer ou mettre à jour des données, mais pas les supprimer.

  • DELETE_ROW_ON_PK : Le framework Flink applique l'opération de suppression basée sur la clé primaire. Pour les opérations de mise à jour, il supprime d'abord les anciennes données puis insère les nouvelles pour garantir l'exactitude des données.

jdbcWriteBatchSize

En mode JDBC, le nombre maximal d'enregistrements à mettre en mémoire tampon dans le sink Hologres avant une écriture par lot.

Integer

Non

256

Unité : lignes.

Remarque

Les paramètres jdbcWriteBatchSize, jdbcWriteBatchByteSize et jdbcWriteFlushInterval ont une relation OU. Si vous définissez les trois paramètres, les données résultantes sont écrites dès qu'une des conditions est remplie.

jdbcWriteBatchByteSize

En mode JDBC, ce paramètre spécifie la taille maximale des données en octets que le sink Hologres met en mémoire tampon avant d'écrire le lot vers la destination.

Long

Non

2 097 152 octets (2 Mo)

Remarque

Les paramètres jdbcWriteBatchSize, jdbcWriteBatchByteSize et jdbcWriteFlushInterval ont une relation OU. Si vous définissez les trois paramètres, les données résultantes sont écrites dès qu'une des conditions est remplie.

jdbcWriteFlushInterval

En mode JDBC, ce paramètre spécifie le temps maximal d'attente du sink Hologres avant d'écrire les données mises en mémoire tampon vers Hologres.

Long

Non

10000

Unité : millisecondes.

Remarque

Les paramètres jdbcWriteBatchSize, jdbcWriteBatchByteSize et jdbcWriteFlushInterval ont une relation OU. Si vous définissez les trois paramètres, les données résultantes sont écrites dès que la condition de l'un d'eux est remplie.

ignoreNullWhenUpdate

Indique s'il faut ignorer les valeurs nulles dans les données entrantes lorsque mutatetype est défini sur INSERT_OR_UPDATE.

Boolean

Non

false

Valeurs valides :

  • false (par défaut) : Écrit les valeurs nulles dans la table sink Hologres.

  • true : Ignore les valeurs nulles dans les données entrantes.

jdbcEnableDefaultForNotNullColumn

Indique s'il faut insérer une valeur par défaut lorsqu'une valeur nulle est écrite dans une colonne NOT NULL qui ne possède pas de valeur par défaut définie.

Boolean

Non

true

Valeurs valides :

  • true (par défaut) : Permet au connecteur d'insérer une valeur par défaut. Les règles sont les suivantes :

    • Pour une colonne de type String, une chaîne vide ("") est écrite.

    • Pour une colonne de type Number, 0 est écrit.

    • Pour les colonnes de type Date, timestamp ou timestamptz, 1970-01-01 00:00:00 est écrit.

  • false : N'insère pas de valeur par défaut. Une exception est levée lorsqu'une valeur nulle est écrite dans une colonne NOT NULL.

remove-u0000-in-text.enabled

Indique s'il faut supprimer le caractère nul (\u0000) des chaînes avant l'écriture.

Boolean

Non

false

Valeurs valides :

  • false (par défaut) : Le connecteur ne traite pas les données. Toutefois, si des données incorrectes sont rencontrées, une opération d'écriture peut lever l'exception suivante : ERROR: invalid byte sequence for encoding "UTF8": 0x00

    Dans ce cas, vous devez soit traiter les données incorrectes dans la table source, soit définir une logique de gestion des données dans votre instruction SQL.

  • true : Le connecteur supprime le caractère \u0000 des chaînes pour éviter les exceptions d'écriture.

deduplication.enabled

Indique s'il faut effectuer une déduplication au sein de chaque lot avant l'écriture en modes jdbc et jdbc_fixed.

Boolean

Non

true

Valeurs valides :

  • true (par défaut) : Active la déduplication au sein d'un lot. Si plusieurs enregistrements ont la même clé primaire, seul le dernier est conservé. Par exemple, considérons des données avec deux champs, où le premier champ est la clé primaire :

    • Les enregistrements INSERT (1,'a') et INSERT (1,'b') arrivent séquentiellement. Après déduplication, seul l'enregistrement ultérieur (1,'b') est écrit dans la table sink Hologres.

    • La table sink Hologres contient déjà l'enregistrement (1,'a'). Si les enregistrements DELETE (1,'a') et INSERT (1,'b') arrivent séquentiellement, seul le dernier enregistrement arrivé (1,'b') est écrit dans Hologres. Ceci est traité comme une mise à jour directe, et non comme une suppression suivie d'une insertion.

  • false : Aucune déduplication n'est effectuée lors du regroupement par lots. Si un nouvel enregistrement a la même clé primaire qu'un enregistrement déjà présent dans le lot actuel, le lot actuel est d'abord écrit. Une fois l'écriture terminée, le nouvel enregistrement est alors traité.

sink.type-normalize-strategy

La stratégie de mappage des types de données.

String

Non

STANDARD

La stratégie utilisée par le sink Hologres pour convertir les types de données amont en types Hologres.

  • STANDARD : Convertit les types Flink CDC en types PostgreSQL (PG) standard.

  • BROADEN : Convertit les types Flink CDC en types Hologres à plage plus large.

  • ONLY_BIGINT_OR_TEXT : Convertit tous les types Flink CDC en BIGINT ou TEXT dans Hologres.

sink.insert.legacy-put-handler

Indique s'il faut utiliser l'ancien gestionnaire Put (Put Handler) pour écrire les données dans Hologres.

Boolean

Non

false

Valeurs valides :

  • false (par défaut) : Le nouveau gestionnaire Put est utilisé pour écrire les données. Le format SQL pour l'opération d'écriture est insert into xxx(c0,c1,...) select unnest(?),unnest(?),... on conflict.

  • true : Écrit les données en utilisant l'ancien gestionnaire Put. Le format SQL pour l'écriture est insert into xxx(c0,c1,...) values (?,?,...),... on conflict; .

table_property.*

Les propriétés de la table physique pour Hologres.

String

Non

Aucune

Lorsque vous créez une table Hologres, vous pouvez définir des propriétés de table physique dans la clause WITH. Des propriétés de table appropriées peuvent aider le système à organiser et interroger les données efficacement.

Avertissement

Le paramètre table_property.distribution_key prend par défaut la valeur de la clé primaire. Ne modifiez pas ce paramètre sauf si vous comprenez parfaitement l'impact, car cela peut affecter l'exactitude des écritures de données.

connection.ssl.mode

Indique s'il faut activer le chiffrement de transport Secure Sockets Layer (SSL) et quel mode utiliser.

String

Non

disable

  • disable (par défaut) : Désactive le chiffrement de transport.

  • require : Active SSL pour chiffrer la liaison de données.

  • verify-ca : Active SSL, chiffre la liaison de données et utilise un certificat CA pour vérifier l'authenticité du serveur Hologres.

  • verify-full : Active SSL, chiffre la liaison de données, utilise un certificat CA pour vérifier l'authenticité du serveur Hologres et vérifie que le nom CN ou DNS dans le certificat correspond au point de terminaison Hologres configuré.

Remarque
  • Hologres V2.1 et versions ultérieures prennent en charge les modes verify-ca et verify-full. Pour plus d'informations, consultez la section Chiffrement de transport.

  • Si vous définissez ce paramètre sur verify-ca ou verify-full, vous devez également configurer le paramètre connection.ssl.root-cert.location.

connection.ssl.root-cert.location

Lorsque le mode de chiffrement de transport nécessite un certificat, spécifie le chemin d'accès au fichier de certificat.

String

Non

Aucune

Lorsque connection.ssl.mode est défini sur verify-ca ou verify-full, vous devez également configurer le chemin d'accès au certificat CA. Vous pouvez télécharger le certificat sur la plateforme à l'aide de la fonctionnalité Gestion des fichiers dans la console Realtime Compute. Une fois le certificat téléchargé, il est stocké dans le répertoire /flink/usrlib. Par exemple, si le fichier de certificat CA est nommé certificate.crt, la valeur du paramètre doit être '/flink/usrlib/certificate.crt'.

Remarque

Pour obtenir un certificat CA, consultez la section Chiffrement de transport - Télécharger un certificat CA.

connection.akv4.enabled

Indique s'il faut activer le mode AKV4 pour se connecter au serveur Hologres.

Boolean

Non

false

N/A.

connection.akv4.region

Lorsque le mode AKV4 est activé, spécifie la région où se trouve le serveur.

String

Non

Aucune

Par exemple, cn-shanghai.

Réutilisation d'un catalogue existant

À partir de VVR 11.5, vous pouvez référencer directement un catalogue Hologres intégré créé sur la page de gestion des données dans une tâche d'ingestion de données Flink CDC. Cela réduit la nécessité de spécifier manuellement les propriétés de connexion.

sink:
  type: hologres
  using.built-in-catalog: my_holo_catalog

Les tâches d'ingestion de données peuvent réutiliser automatiquement les paramètres de catalogue Hologres suivants :

  • endpoint

  • username

  • password

  • dbname

Pour remplacer ces paramètres réutilisés automatiquement, spécifiez explicitement les paramètres YAML correspondants, qui ont une priorité plus élevée.

Mappage des types de données

Utilisez le paramètre sink.type-normalize-strategy pour définir la stratégie de conversion des données amont vers les types Hologres.

Remarque
  • Activez sink.type-normalize-strategy lors du premier démarrage d'une tâche YAML. Si vous l'activez après le démarrage de la tâche, supprimez la table aval et redémarrez la tâche sans état pour que le paramètre prenne effet.

  • Actuellement, les types de tableaux ne prennent en charge que INTEGER, BIGINT, FLOAT, DOUBLE, BOOLEAN, CHAR et VARCHAR.

  • Hologres ne prend pas en charge le type numeric comme clé primaire. Si le type d'une clé primaire correspond à numeric, le système le convertit en type varchar.

STANDARD

Lorsque sink.type-normalize-strategy est défini sur STANDARD, les types sont mappés comme suit :

Type Flink CDC

Type Hologres

CHAR

bpchar

STRING

text

VARCHAR

text (si length > 10485760)

varchar (si length <= 10485760)

BOOLEAN

bool

BINARY

bytea

VARBINARY

DECIMAL

numeric

TINYINT

int2

SMALLINT

INTEGER

int4

BIGINT

int8

FLOAT

float4

DOUBLE

float8

DATE

date

TIME_WITHOUT_TIME_ZONE

time

TIMESTAMP_WITHOUT_TIME_ZONE

timestamp

TIMESTAMP_WITH_LOCAL_TIME_ZONE

timestamptz

ARRAY

Tableaux des types d'éléments correspondants

MAP

Non pris en charge

ROW

Non pris en charge

BROADEN

Lorsque sink.type-normalize-strategy est défini sur BROADEN, les types Flink CDC sont convertis en types Hologres à plage plus large. Les types sont mappés comme suit :

Type Flink CDC

Type Hologres

CHAR

text

STRING

VARCHAR

BOOLEAN

bool

BINARY

bytea

VARBINARY

DECIMAL

numeric

TINYINT

int8

SMALLINT

INTEGER

BIGINT

FLOAT

float8

DOUBLE

DATE

date

TIME_WITHOUT_TIME_ZONE

time

TIMESTAMP_WITHOUT_TIME_ZONE

timestamp

TIMESTAMP_WITH_LOCAL_TIME_ZONE

timestamptz

ARRAY

Tableaux des types d'éléments correspondants

MAP

Non pris en charge

ROW

Non pris en charge

ONLY_BIGINT_OR_TEXT

Lorsque sink.type-normalize-strategy est défini sur ONLY_BIGINT_OR_TEXT, tous les types Flink CDC sont convertis en types BIGINT ou STRING dans Hologres. Le mappage des types de données est le suivant :

Type Flink CDC

Type Hologres

TINYINT

int8

SMALLINT

INTEGER

BIGINT

BOOLEAN

text

BINARY

VARBINARY

DECIMAL

FLOAT

DOUBLE

DATE

TIME_WITHOUT_TIME_ZONE

TIMESTAMP_WITHOUT_TIME_ZONE

TIMESTAMP_WITH_LOCAL_TIME_ZONE

ARRAY

Tableaux des types d'éléments correspondants

MAP

Non pris en charge

ROW

Non pris en charge

Écriture dans des tables partitionnées

Vous pouvez utiliser le sink Hologres avec une transformation pour écrire les données amont dans une table partitionnée Hologres.

  • La partition key doit faire partie de la primary key. Si vous utilisez une colonne non clé primaire des données amont comme partition key, les clés primaires des tables amont et aval peuvent devenir incohérentes, entraînant des divergences de données lors de la synchronisation.

  • Hologres prend en charge les colonnes de types de données TEXT, VARCHAR et INT comme partition key. Depuis la version 1.3.22, les colonnes de type de données DATE sont également prises en charge.

  • Pour créer automatiquement des tables partitionnées enfants, définissez le paramètre createparttable sur true. Sinon, vous devez les créer manuellement.

Consultez la section Écriture de données dans une table partitionnée pour un exemple.

Synchronisation du schéma de table

Un pipeline YAML CDC utilise différentes stratégies pour gérer les modifications de schéma de table, configurables via le paramètre de niveau pipeline schema.change.behavior. Les valeurs valides pour schema.change.behavior sont IGNORE, LENIENT, TRY_EVOLVE, EVOLVE et EXCEPTION. Hologres Sink ne prend pas actuellement en charge la stratégie TRY_EVOLVE. Les stratégies LENIENT et EVOLVE impliquent des modifications de schéma de table. Les sections suivantes décrivent comment ces deux modes gèrent différents événements de modification de schéma.

LENIENT (par défaut)

En mode LENIENT, les modifications de schéma sont gérées comme suit :

  • Ajout d'une colonne nullable : La colonne correspondante est automatiquement ajoutée à la fin de la table sink et ses données sont synchronisées.

  • Suppression d'une colonne nullable : La colonne n'est pas supprimée de la table sink. À la place, la colonne est automatiquement renseignée avec des valeurs NULL.

  • Ajout d'une colonne non nullable : Une colonne nullable correspondante est automatiquement ajoutée à la fin de la table sink et ses données sont synchronisées. Pour les lignes existantes, cette nouvelle colonne est automatiquement renseignée avec NULL.

  • Renommage d'une colonne : Cette opération est traitée comme la suppression d'une colonne et l'ajout d'une nouvelle. Une nouvelle colonne portant le nom spécifié est ajoutée à la fin de la table sink, et la colonne portant le nom d'origine est automatiquement renseignée avec des valeurs NULL. Par exemple, si col_a est renommé en col_b, une colonne col_b est ajoutée à la fin de la table sink, et la colonne col_a est automatiquement renseignée avec des valeurs NULL.

  • Modification du type de données d'une colonne : Non pris en charge. Comme Hologres ne prend pas en charge la modification du type de données d'une colonne, vous devez utiliser le paramètre sink.type-normalize-strategy.

  • Les modifications de schéma suivantes ne sont pas prises en charge :

    • Modifications des contraintes, telles que la clé primaire ou l'index.

    • Suppression d'une colonne non nullable.

    • Passage d'une colonne de NOT NULL à NULLABLE.

EVOLVE

En mode EVOLVE, les modifications de schéma sont gérées comme suit :

  • Ajout d'une colonne nullable : Pris en charge.

  • Suppression d'une colonne nullable : Non pris en charge.

  • Ajout d'une colonne non nullable : Une nouvelle colonne nullable est ajoutée à la table sink.

  • Renommage d'une colonne : Pris en charge. La colonne d'origine est renommée dans la table sink.

  • Modification du type de données d'une colonne : Non pris en charge. Comme Hologres ne prend pas en charge la modification du type de données d'une colonne, vous devez utiliser le paramètre sink.type-normalize-strategy.

  • Les modifications de schéma suivantes ne sont pas prises en charge :

    • Modifications des contraintes, telles que la clé primaire ou l'index.

    • Suppression d'une colonne non nullable.

    • Passage d'une colonne de NOT NULL à NULLABLE.

Avertissement

En mode EVOLVE, effectuer un redémarrage sans état sans supprimer la table sink peut entraîner l'échec du pipeline en raison d'incohérences de schéma entre les données amont et la table sink. Vous devez alors ajuster manuellement le schéma de la table sink.

Consultez la section Activation du mode EVOLVE pour un exemple.

Exemples de code

Élargissement des types

Utilisez le paramètre sink.type-normalize-strategy pour configurer l'élargissement des types.

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.test_source_table
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # Map CDC data types to broader Hologres types.
  sink.type-normalize-strategy: BROADEN

pipeline:
  name: MySQL to Hologres Pipeline

Écriture dans une table partitionnée

Convertit le champ d'horodatage create_time en type date et l'utilise comme clé de partition pour la table Hologres.

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.test_source_table
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # Automatically create partitioned tables if they do not exist.
  createparttable: true
 
transform:
  - source-table: test_db.test_source_table
    projection: \*, DATE_FORMAT(CAST(create_time AS TIMESTAMP), 'yyyy-MM-dd') as partition_key
    primary-keys: id, create_time, partition_key
    partition-keys: partition_key
    description: add partition key 

pipeline:
  name: MySQL to Hologres Pipeline

Activation du mode EVOLVE****

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.test_source_table
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # Automatically create partitioned tables if they do not exist.
  createparttable: true

pipeline:
  name: MySQL to Hologres Pipeline
  schema.change.behavior: evolve

Synchronisation de table unique

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.test_source_table
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # Map CDC data types to broader Hologres types.
  sink.type-normalize-strategy: BROADEN

pipeline:
  name: MySQL to Hologres Pipeline

Synchronisation complète de la base de données

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.\.*
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # Map CDC data types to broader Hologres types.
  sink.type-normalize-strategy: BROADEN

pipeline:
  name: MySQL to Hologres Pipeline

Fusion de tables fragmentées

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.user\.*
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # Map CDC data types to broader Hologres types.  
  sink.type-normalize-strategy: BROADEN
  
route:
  # All sharded tables in the MySQL test_db database are merged into a single Hologres table named test_db.user.
  - source-table: test_db.user\.*
    sink-table: test_db.user

pipeline:
  name: MySQL to Hologres Pipeline

Synchronisation vers un schéma spécifié

Dans Hologres, un schéma correspond à une base de données dans MySQL. Vous pouvez spécifier le schéma pour les tables sink.

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.user\.*
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # Map CDC data types to broader Hologres types.
  sink.type-normalize-strategy: BROADEN
  
route:
  # Synchronize all tables from the MySQL test_db database to the Hologres test_db2 schema while keeping the original table names.
  - source-table: test_db.\.*
    sink-table: test_db2.<>
    replace-symbol: <>

pipeline:
  name: MySQL to Hologres Pipeline

Synchronisation de nouvelles tables sans redémarrage

Pour synchroniser en temps réel les tables nouvellement ajoutées pendant l'exécution d'une tâche, définissez scan.binlog.newly-added-table.enable = true.

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.\.*
  server-id: 5401-5499
  # Automatically captures new tables created while the job is running.  
  scan.binlog.newly-added-table.enabled: true

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # Map CDC data types to broader Hologres types.
  sink.type-normalize-strategy: BROADEN

pipeline:
  name: MySQL to Hologres Pipeline

Ajout de tables existantes au redémarrage

Pour inclure une table existante dans la synchronisation, définissez scan.newly-added-table.enabled sur true et redémarrez la tâche.

Avertissement

N'utilisez pas scan.newly-added-table.enabled sur une tâche qui a été précédemment exécutée avec scan.binlog.newly-added-table.enabled. Cette combinaison provoque une duplication des données au redémarrage.

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.\.*
  server-id: 5401-5499
  scan.startup.mode: initial
  # On restart, the job scans for new tables matching the `tables` parameter and performs a snapshot.
  # Note: This parameter must be used with scan.startup.mode: initial.
  scan.newly-added-table.enabled: true

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # Map CDC data types to broader Hologres types.
  sink.type-normalize-strategy: BROADEN

pipeline:
  name: MySQL to Hologres Pipeline

Exclusion de tables

source:
  type: mysql
  name: MySQL Source
  hostname: <yourHostname>
  port: 3306
  username: flink
  password: ${secret_values.password}
  tables: test_db.\.*
  # Excludes tables that match this regular expression.
  tables.exclude: test_db.table1
  server-id: 5401-5499

sink:
  type: hologres
  name: Hologres Sink
  endpoint: <yourEndpoint>
  dbname: <yourDbname>
  username: ${secret_values.ak_id}
  password: ${secret_values.ak_secret}
  # Map CDC data types to broader Hologres types.
  sink.type-normalize-strategy: BROADEN

pipeline:
  name: MySQL to Hologres Pipeline

Documents connexes