Este tópico descreve os módulos source e sink para jobs de ingestão de dados do Flink Change Data Capture (CDC) e lista os conectores compatíveis.
Conectores compatíveis
Conector | Tipo compatível | |
Source | Sink | |
Nota Compatível com conexões ao ApsaraDB RDS for MySQL, PolarDB for MySQL e MySQL autogerenciado. | √ | × |
× | √ | |
√ Nota Compatível apenas no VVR 8.0.10 ou superior. | √ | |
× | √ | |
× | √ | |
× | √ | |
√ Nota Compatível apenas no VVR 11.1 ou superior. | × | |
√ Nota Compatível apenas no VVR 11.2 ou superior. | × | |
× | √ Nota Compatível apenas no VVR 11.1 ou superior. | |
× | √ Nota Compatível apenas no VVR 11.1 ou superior. | |
√ Nota Compatível apenas no VVR 11.4 ou superior. | × | |
× | √ | |
× | √ Nota Compatível apenas no VVR 11.6 ou superior. | |
Configuração do conector
Configure parâmetros para conectores source e sink em um job de ingestão de dados do Flink CDC. Para obter detalhes sobre os conectores compatíveis e seus parâmetros, consulte as seções a seguir.
# Source module
source:
type: mysql # Or another connector identifier
name: MySQL Source
# Other parameters. Use key: value pairs.
# Sink module
sink:
type: paimon # Or another connector identifier
name: Paimon Sink
# Other parameters. Use key: value pairs.
Configuração geral
Parâmetro | Descrição | Obrigatório | Tipo | Padrão | Observações |
type | Tipo de conector para o source ou sink. | Sim | String | Nenhum | Nenhum |
name | Nome do nó. | Não | String | Nenhum | Nenhum |
using.built-in-catalog | Reutiliza um catálogo integrado. | Não | String | Nenhum |
|
Reutilizar informações de conexão de um catálogo
A partir do VVR 11.5, referencie diretamente um catálogo integrado criado na página Data Management para seu job de ingestão de dados do Flink CDC. Essa funcionalidade permite reutilizar propriedades de conexão, como URL, nome de usuário e senha, simplificando a configuração manual.
Sintaxe
source:
type: mysql
using.built-in-catalog: mysql_rds_catalog
sink:
type: paimon
using.built-in-catalog: paimon_dlf_catalog
Use o parâmetro using.built-in-catalog nos módulos source e sink para referenciar um catálogo integrado existente.
No exemplo anterior, os metadados do catálogo para mysql_rds_catalog já incluem parâmetros obrigatórios como hostname, username e password. Portanto, não é necessário fornecer esses parâmetros novamente no job yaml.
Notas de uso
Os conectores a seguir oferecem suporte à reutilização de informações de conexão de um catálogo:
MySQL (source)
Kafka (source)
Upsert Kafka (sink)
StarRocks (sink)
Hologres (sink)
Paimon (sink)
Simple Log Service (source)
Iceberg (sink)
O sistema ignora parâmetros de catálogo incompatíveis com a configuração yaml do CDC. Para mais informações, consulte a lista de parâmetros de cada conector.
Configuração do source
Parâmetro | Descrição | Obrigatório | Tipo | Padrão | Observações |
source-expand | Define a estratégia de distribuição para os dados emitidos pelo source. | Não | Consulte a seção de sintaxe para detalhes de configuração. | Nenhum |
|
source-expand
O parâmetro source-expand distribui ou replica dados antes do processamento pelos módulos Transform e Route.
Sintaxe
source:
type: mysql
host: localhost
port: 3306
username: admin
password: pass
tables: mydb.orders
source-expand:
# Replicate mydb.orders into three tables, each processed differently before writing to the sink.
- input-table: mydb.orders
output-table: [ dwd.orders_full, dws.orders_summary, ads.orders_report ]
# Apply different transformations to each of the three expanded tables
transform:
# DWD layer: Retain all fields, add a calculated column
- source-table: dwd.orders_full
projection: "*, amount * discount as final_price"
# DWS layer: Retain only key fields, filter out small orders
- source-table: dws.orders_summary
projection: order_id, user_id, amount, order_status
filter: amount > 100
primary-keys: order_id
# ADS layer: Retain only fields needed for reporting, filter out canceled orders
- source-table: ads.orders_report
projection: order_id, user_id, amount, TO_UPPER(order_status) as status
filter: order_status <> 'CANCELLED'
primary-keys: order_id
# Route the three expanded tables to different sink tables
route:
- source-table: dwd.orders_full
sink-table: starrocks_dwd.orders_full_detail
- source-table: dws.orders_summary
sink-table: starrocks_dws.orders_summary
- source-table: ads.orders_report
sink-table: starrocks_ads.orders_report
sink:
type: starrocks
name: sink-starrocks
jdbc-url: jdbc:mysql://localhost:9030
load-url: localhost:8030
username: root
password: pass
Conforme o exemplo, essa configuração replica dados da tabela mydb.orders para três tabelas lógicas: dwd.orders_full, dws.orders_summary e ads.orders_report. O sistema processa cada tabela de forma diferente antes de gravá-la em uma tabela sink separada.
Notas de uso
Disponível no VVR 11.6 ou superior.
-
Por padrão, o sistema não mantém a tabela original após a distribuição dos dados. No exemplo de sintaxe, a tabela
mydb.ordersé removida do fluxo de dados. Para manter a tabela original, adicione-a à listaoutput-table.source-expand: - input-table: mydb.orders output-table: [ mydb.orders, db1.orders, db2.orders ] Os parâmetros
input-tableeoutput-tablenão aceitam expressões regulares.
Configuração do sink
|
Parâmetro |
Descrição |
Obrigatório |
Tipo |
Padrão |
Observações |
|
include.schema.changes |
Define quais tipos de alterações de esquema aplicar. |
Não |
List<String> |
Nenhum |
Por padrão, todas as alterações de esquema são sincronizadas. |
|
exclude.schema.changes |
Define quais tipos de alterações de esquema excluir. |
Não |
List<String> |
Nenhum |
Tem prioridade maior que |
Para mais informações sobre como usar include.schema.changes e exclude.schema.changes, consulte Configuração de sincronização de alterações de esquema.