O Realtime Compute for Apache Flink oferece um recurso avançado de ingestão de dados baseado no Flink CDC. Este guia demonstra como crie um job de ingestão de dados do Flink CDC para sincronizar um banco de dados MySQL inteiro com um banco de dados StarRocks.
Pré-requisitos
Um workspace do Flink criado. Para mais informações, consulte Ativar o Realtime Compute for Apache Flink.
-
Armazenamentos de dados de source e sink
Uma instância do ApsaraDB RDS for MySQL criada. Para mais informações, consulte (Descontinuado, redirecionado para "Etapa 1") Criar rapidamente uma instância do ApsaraDB RDS for MySQL.
Uma instância do EMR Serverless StarRocks criada. Para mais informações, consulte Procedimento.
NotaA instância do ApsaraDB RDS for MySQL e a instância do EMR Serverless StarRocks devem estar na mesma Virtual Private Cloud (VPC) do seu workspace do Flink. Caso estejam em VPCs diferentes, estabeleça uma conexão de rede e configure uma lista de permissões para a instância do ApsaraDB RDS for MySQL. Para mais informações, consulte Como acesso outros serviços entre VPCs?, Como acesso a Internet? e Como configuro uma lista de permissões?.
Informações de fundo
Considere que sua instância do ApsaraDB RDS for MySQL possua um banco de dados chamado order_dw_mysql contendo três tabelas de negócios: orders, orders_pay e product_catalog. Para sincronizar essas tabelas e seus dados com o banco de dados order_dw_sr no StarRocks, siga estas etapas:
Etapa 1: Preparar dados de teste do MySQL
-
Crie um banco de dados e uma conta.
Crie um banco de dados chamado order_dw_mysql e uma conta padrão com permissões de leitura e gravação nele. Para mais informações, consulte (Descontinuado, redirecionado para "Etapa 1") Criar um banco de dados e uma conta e Gerenciar bancos de dados.
-
Faça login na instância do ApsaraDB RDS for MySQL usando o Data Management (DMS).
Para mais informações, consulte (Descontinuado, redirecionado para "Etapa 2") Fazer login em uma instância do ApsaraDB RDS for MySQL usando o DMS.
-
Na janela do SQL Console, insira os comandos abaixo e clique em Execute para criar três tabelas de negócios e inserir dados.
CREATE TABLE `orders` ( order_id bigint not null primary key, user_id varchar(50) not null, shop_id bigint not null, product_id bigint not null, buy_fee numeric(20,2) not null, create_time timestamp not null, update_time timestamp not null default now(), state int not null ); CREATE TABLE `orders_pay` ( pay_id bigint not null primary key, order_id bigint not null, pay_platform int not null, create_time timestamp not null ); CREATE TABLE `product_catalog` ( product_id bigint not null primary key, catalog_name varchar(50) not null ); -- Prepare data INSERT INTO product_catalog VALUES(1, 'phone_aaa'),(2, 'phone_bbb'),(3, 'phone_ccc'),(4, 'phone_ddd'),(5, 'phone_eee'); INSERT INTO orders VALUES (100001, 'user_001', 12345, 1, 5000.05, '2023-02-15 16:40:56', '2023-02-15 18:42:56', 1), (100002, 'user_002', 12346, 2, 4000.04, '2023-02-15 15:40:56', '2023-02-15 18:42:56', 1), (100003, 'user_003', 12347, 3, 3000.03, '2023-02-15 14:40:56', '2023-02-15 18:42:56', 1), (100004, 'user_001', 12347, 4, 2000.02, '2023-02-15 13:40:56', '2023-02-15 18:42:56', 1), (100005, 'user_002', 12348, 5, 1000.01, '2023-02-15 12:40:56', '2023-02-15 18:42:56', 1), (100006, 'user_001', 12348, 1, 1000.01, '2023-02-15 11:40:56', '2023-02-15 18:42:56', 1), (100007, 'user_003', 12347, 4, 2000.02, '2023-02-15 10:40:56', '2023-02-15 18:42:56', 1); INSERT INTO orders_pay VALUES (2001, 100001, 1, '2023-02-15 17:40:56'), (2002, 100002, 1, '2023-02-15 17:40:56'), (2003, 100003, 0, '2023-02-15 17:40:56'), (2004, 100004, 0, '2023-02-15 17:40:56'), (2005, 100005, 0, '2023-02-15 18:40:56'), (2006, 100006, 0, '2023-02-15 18:40:56'), (2007, 100007, 0, '2023-02-15 18:40:56');
Etapa 2: Desenvolver um job do Flink CDC
Faça login no console de gerenciamento do Realtime Compute for Apache Flink.
Clique em Console para acessar o workspace do projeto.
No painel de navegação à esquerda, escolha .
Clique no ícone
, clique em New Draft with Template, selecione MySQL to StarRocks data synchronization e clique em Next.Insira um Job Name e um Storage Location, selecione uma Engine Version e clique em OK.
-
Configure o código YAML do job.
O código a seguir apresenta um exemplo de sincronização de todas as tabelas do banco de dados
order_dw_mysqlno MySQL para o banco de dadosorder_dw_srno StarRocks.source: type: mysql hostname: rm-bp1rk934iidc3****.mysql.rds.aliyuncs.com port: 3306 username: ${secret_values.mysqlusername} password: ${secret_values.mysqlpassword} tables: order_dw_mysql.\.* server-id: 8601-8604 # (Optional) Synchronize data from tables that are newly created during the incremental phase. scan.binlog.newly-added-table.enabled: true # (Optional) Synchronize table and column comments. include-comments.enabled: true # (Optional) Prioritize the distribution of unbounded chunks to prevent potential TaskManager OutOfMemory errors. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Optional) Enable parsing filters to accelerate reading. scan.only.deserialize.captured.tables.changelog.enabled: true sink: type: starrocks name: StarRocks Sink jdbc-url: jdbc:mysql://fe-c-b76b6aa51807****-internal.starrocks.aliyuncs.com:9030 load-url: fe-c-b76b6aa51807****-internal.starrocks.aliyuncs.com:8030 username: ${secret_values.starrocksusername} password: ${secret_values.starrockspassword} table.create.properties.replication_num: 1 sink.buffer-flush.interval-ms: 5000 # Flush data every 5 seconds. route: - source-table: order_dw_mysql.\.* sink-table: order_dw_sr.<> replace-symbol: <> description: route all tables in source_db to sink_db pipeline: name: Sync MySQL Database to StarRocksA tabela a seguir descreve os parâmetros de configuração necessários para este exemplo. Para mais detalhes sobre os parâmetros de ingestão de dados, consulte MySQL e StarRocks.
NotaJobs YAML suportam apenas variáveis de projeto. Utilize variáveis para evitar que informações como senhas sejam exibidas em texto simples. Para mais informações, consulte Gerenciamento de Variáveis.
Categoria
Parâmetro
Descrição
Valor de exemplo
source
hostname
Endereço IP ou nome do host do banco de dados MySQL.
Recomendamos o uso do endpoint interno.
rm-bp1rk934iidc3****.mysql.rds.aliyuncs.comport
Número da porta do serviço de banco de dados MySQL.
3306
username
Nome de usuário e senha do banco de dados MySQL. Utilize as credenciais da conta criada na Etapa 1: Preparar dados de teste no ApsaraDB RDS for MySQL.
${secret_values.mysqlusername}password
${secret_values.mysqlpassword}tables
Nomes das tabelas do MySQL. É possível usar expressões regulares para ler dados de várias tabelas.
Neste tópico, todas as tabelas e dados do banco de dados order_dw_mysql são sincronizados.
order_dw_mysql.\.*
server-id
ID numérico exclusivo para a conexão do cliente ao banco de dados.
8601-8604
sink
jdbc-url
URL de conexão JDBC.
Especifique o endereço IP e a porta de consulta do Frontend (FE) no formato
jdbc:mysql://ip:port.Na aba Instance Details do console E-MapReduce, visualize o internal endpoint e a query port do FE da instância de destino.
jdbc:mysql://fe-c-b76b6aa51807****-internal.starrocks.aliyuncs.com:9030load-url
URL do serviço HTTP usada para conectar ao nó FE.
Na aba Instance Details do console E-MapReduce, visualize o internal endpoint e a HTTP port do FE da instância de destino.
fe-c-b76b6aa51807****-internal.starrocks.aliyuncs.com:8030username
Credenciais para conexão com o StarRocks.
Utilize as credenciais configuradas durante a criação da instância do StarRocks.
NotaEste exemplo utiliza variáveis para evitar a exposição de credenciais em texto simples. Para mais informações, consulte Gerenciar variáveis.
${secret_values.starrocksusername}password
${secret_values.starrockspassword}sink.buffer-flush.interval-ms
Intervalo de liberação (flush) do buffer interno.
Um intervalo curto (5 segundos) é utilizado porque este exemplo contém poucos dados, permitindo visualizar os resultados rapidamente.
5000
route
source-table
Tabela ou tabelas de origem a serem roteadas.
Use uma expressão regular para corresponder a várias tabelas. Por exemplo,
order_dw_mysql.\.*roteia todas as tabelas do banco de dadosorder_dw_mysql.order_dw_mysql.\.*
sink-table
Padrão da tabela de destino para os dados roteados.
Utilize o símbolo definido no parâmetro
replace-symbolcomo espaço reservado para cada nome de tabela de origem, possibilitando o roteamento muitos-para-muitos.Para mais informações sobre regras de roteamento, consulte Módulo Route.
order_dw_sr.<>
replace-symbol
Espaço reservado para o nome da tabela de origem usado na correspondência de padrões.
<>
Clique em Deploy.
Etapa 3: Iniciar o job do Flink CDC
Na página Data Ingestion, clique em Deploy e, em seguida, clique em OK na caixa de diálogo exibida.
Na página , localize o job YAML desejado e clique em Start na coluna Actions.
-
Clique em Start.
Neste exemplo, selecione Initial Mode. Para mais informações sobre os parâmetros, consulte Iniciar um job. Após o início do job, monitore seu status na página Deployments.
Etapa 4: Verificar resultados no StarRocks
Depois que o job entrar no estado RUNNING, verifique os dados no StarRocks.
Conectar a uma instância do EMR Serverless StarRocks usando o EMR StarRocks Manager.
-
No painel de navegação à esquerda, clique em SQL Editor. Na aba Database, clique no ícone de atualização
.Um banco de dados chamado order_dw_sr aparecerá sob default_catalog.
-
Na aba Query List, clique em + File para criar um Query Script. Insira as instruções SQL a seguir e clique em Run.
SELECT * FROM default_catalog.order_dw_sr.orders order by order_id; SELECT * FROM default_catalog.order_dw_sr.orders_pay order by pay_id; SELECT * FROM default_catalog.order_dw_sr.product_catalog order by product_id; -
Visualize os resultados abaixo dos comandos.
Os resultados mostram que as tabelas e os dados do banco de dados MySQL agora existem no StarRocks.
As tabelas sincronizadas incluem
default_catalog.order_dw_sr.orders,default_catalog.order_dw_sr.orders_payedefault_catalog.order_dw_sr.product_catalog. Execute instruções SELECT para consultar cada tabela e verificar a integridade dos dados.
Documentação relacionada
Para obter etapas detalhadas sobre o desenvolvimento de um job de ingestão de dados do Flink CDC, consulte Desenvolver um job de ingestão de dados do Flink CDC.
Para mais informações sobre os módulos source, sink, transform e route para jobs de ingestão de dados do Flink CDC, consulte Referência de desenvolvimento de job de ingestão de dados do Flink CDC.