Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Módulos source e sink do Flink CDC

Última atualização: Jun 27, 2026

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

MySQL

Nota

Compatível com conexões ao ApsaraDB RDS for MySQL, PolarDB for MySQL e MySQL autogerenciado.

×

Paimon

×

Kafka

Nota

Compatível apenas no VVR 8.0.10 ou superior.

Upsert Kafka

×

StarRocks

×

Hologres

×

Simple Log Service

Nota

Compatível apenas no VVR 11.1 ou superior.

×

MongoDB

Nota

Compatível apenas no VVR 11.2 ou superior.

×

MaxCompute

×

Nota

Compatível apenas no VVR 11.1 ou superior.

SelectDB

×

Nota

Compatível apenas no VVR 11.1 ou superior.

PostgreSQL CDC (Public Preview)

Nota

Compatível apenas no VVR 11.4 ou superior.

×

Print

×

Iceberg

×

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)

Nota

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

  • Disponível no VVR 11.6 ou superior.

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 à lista output-table.

    source-expand:
      - input-table: mydb.orders
        output-table: [ mydb.orders, db1.orders, db2.orders ]
  • Os parâmetros input-table e output-table nã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 include.schema.changes.

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.