Saiba como construir um streaming data lakehouse usando Realtime Compute for Apache Flink, Apache Paimon e StarRocks.
Contexto
Com a crescente digitalização dos negócios, a demanda por dados em tempo real aumenta continuamente. O método tradicional de construção de um data warehouse offline envolve o agendamento de jobs offline para mesclar periodicamente novos dados em estruturas em camadas, como ODS, DWD, DWS e ADS. No entanto, essa abordagem apresenta dois problemas significativos: alta latência e alto custo. Jobs offline geralmente são agendados de hora em hora ou até diariamente, de modo que os consumidores downstream só conseguem visualizar dados com pelo menos uma hora ou um dia de atraso. Além disso, as atualizações frequentemente exigem a sobrescrita de partições inteiras. Esse processo ineficiente relê todos os dados originais de uma partição para mesclar as novas alterações.
Construir um streaming data lakehouse com Realtime Compute for Apache Flink e Apache Paimon resolve os problemas de latência e custo dos data warehouses offline tradicionais. As capacidades de processamento em tempo real do Flink permitem que os dados fluam continuamente entre as camadas do warehouse. Enquanto isso, as capacidades eficientes de atualização do Paimon entregam alterações de dados aos consumidores downstream com latência de apenas alguns minutos. Como resultado, o streaming data lakehouse oferece vantagens significativas tanto em baixa latência quanto em custo-benefício.
Para obter mais informações sobre o Apache Paimon, consulte seus Recursos e o site oficial do Apache Paimon.
Arquitetura e benefícios
Arquitetura
O Realtime Compute for Apache Flink é um mecanismo poderoso de processamento de fluxo que processa eficientemente grandes volumes de dados em tempo real. O Paimon é um formato de armazenamento de lake unificado para processamento em fluxo e em lote, compatível com atualizações de alto throughput e consultas de baixa latência. O Paimon integra-se profundamente ao Flink para fornecer uma solução completa de streaming lakehouse. A seguir, veja a arquitetura para construir um streaming lakehouse usando Flink e Paimon.
O Flink grava dados das fontes de dados no Paimon, criando a camada ODS.
O Flink assina o changelog da camada ODS, transforma os dados e os grava de volta no Paimon como a camada DWD.
O Flink assina o changelog da camada DWD, transforma os dados e os grava de volta no Paimon como a camada DWS.
Por fim, o StarRocks no E-MapReduce lê a tabela externa do Paimon para atender às consultas da aplicação.

Benefícios
Esta solução oferece os seguintes benefícios:
Cada camada de dados no Paimon propaga alterações para sistemas downstream em questão de minutos. Isso reduz a latência de um data warehouse offline tradicional de horas ou até dias para minutos.
Cada camada de dados no Paimon ingere diretamente dados de alteração sem sobrescrever partições. Isso reduz significativamente os custos de atualização e correção de dados em um data warehouse offline tradicional e resolve os desafios de consulta, atualização e correção de dados nas camadas intermediárias.
A solução possui um modelo unificado e uma arquitetura simplificada. A lógica do pipeline de ETL é implementada usando Flink SQL. Os dados nas camadas ODS, DWD e DWS são armazenados no Paimon. Isso reduz a complexidade arquitetônica e melhora a eficiência do processamento de dados.
Esta solução baseia-se em três capacidades principais do Paimon, conforme detalhado na tabela a seguir.
|
Capacidade principal |
Descrição |
|
Atualizações de tabela de chave primária |
O Paimon usa internamente uma estrutura de dados LSM tree para permitir atualizações eficientes de dados. Para obter mais informações sobre tabelas de chave primária do Paimon e suas estruturas de dados subjacentes, consulte Primary Key Table e File Layouts. |
|
Changelog Producer |
O Paimon pode gerar dados incrementais completos para qualquer fluxo de dados de entrada, onde cada registro |
|
Merge Engine |
Quando uma tabela de chave primária do Paimon recebe vários registros com a mesma chave primária, o Merge Engine mescla esses registros em um único registro para manter a unicidade da chave. O Paimon suporta uma ampla variedade de comportamentos de mesclagem, como |
Casos de uso
Este artigo usa uma plataforma de e-commerce como exemplo para demonstrar como construir um streaming data lakehouse. Esta solução processa e limpa dados para suportar consultas de aplicações upstream. Ao utilizar camadas de dados, esta solução permite a reutilização de dados em múltiplos cenários de negócios, como painéis de transações, análise de dados comportamentais, tags de perfil de usuário e recomendações personalizadas.

Construa a camada ODS: Ingestão em tempo real do banco de dados de negócios. O Flink ingere três tabelas de negócios do MySQL,
orders,orders_payeproduct_catalog, no OSS em tempo real. Esses dados, armazenados no formato Paimon, constituem a camada ODS.Construa a camada DWD: Tabela larga temática. As tabelas
orders,product_catalogeorders_paysão unidas usando o mecanismo de mesclagem de atualização parcial do Paimon, criando uma tabela larga temática para a camada DWD e um changelog com latência de minutos.Construa a camada DWS: Cálculo de métricas. O Flink consome o changelog da tabela larga em tempo real. Ele usa o mecanismo de mesclagem de pré-agregação do Paimon para criar uma tabela de agregação intermediária,
dwm_users_shops. Esse processo gera finalmente duas tabelas para a camada DWS:dws_userspara métricas de agregação de usuários edws_shopspara métricas de agregação de lojas.
Pré-requisitos
Ative o Data Lake Formation (DLF). Recomendamos usar o DLF 2.5 como serviço de armazenamento. Para obter mais informações, consulte Introdução ao DLF.
Crie um workspace do Realtime Compute for Apache Flink. Para obter mais informações, consulte Criar um workspace.
Crie uma instância Serverless StarRocks. Para obter mais informações, consulte Introdução a uma instância Serverless StarRocks.
Sua instância StarRocks, o DLF e o workspace do Flink devem estar na mesma região.
Limites
Esta solução de streaming lakehouse requer Ververica Runtime (VVR) 11.1.0 ou posterior.
Construir um streaming data lakehouse
Preparar a fonte de dados MySQL CDC****
Este tutorial usa o ApsaraDB RDS for MySQL como exemplo. Você criará um banco de dados chamado order_dw e três tabelas com dados de amostra.
-
(Descontinuado. Redireciona para a Etapa 1.) Criar uma instância ApsaraDB RDS for MySQL.
ImportanteCertifique-se de que a instância ApsaraDB RDS for MySQL e seu workspace do Realtime Compute for Apache Flink estejam na mesma VPC. Se estiverem em VPCs diferentes, consulte Como acesso outros serviços entre VPCs?
-
(Descontinuado. Redireciona para a Etapa 1.) Criar um banco de dados e uma conta.
Crie um banco de dados chamado order_dw. Em seguida, crie uma conta privilegiada ou uma conta padrão com permissões de leitura e gravação para o banco de dados order_dw.
Crie três tabelas e insira dados nelas.
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 bigint 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, '2023-02-15 16:40:56', '2023-02-15 18:42:56', 1), (100002, 'user_002', 12346, 2, 4000, '2023-02-15 15:40:56', '2023-02-15 18:42:56', 1), (100003, 'user_003', 12347, 3, 3000, '2023-02-15 14:40:56', '2023-02-15 18:42:56', 1), (100004, 'user_001', 12347, 4, 2000, '2023-02-15 13:40:56', '2023-02-15 18:42:56', 1), (100005, 'user_002', 12348, 5, 1000, '2023-02-15 12:40:56', '2023-02-15 18:42:56', 1), (100006, 'user_001', 12348, 1, 1000, '2023-02-15 11:40:56', '2023-02-15 18:42:56', 1), (100007, 'user_003', 12347, 4, 2000, '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');
Gerenciar metadados
Criar um catálogo Paimon
Faça logon no console do Realtime Compute for Apache Flink.
No painel de navegação à esquerda, acesse Metadata Management e clique em Create Catalog.
Na aba Built-in Catalog, clique em Apache Paimon e, em seguida, clique em Next.
-
Configure os parâmetros a seguir, selecione DLF como tipo de armazenamento e clique em OK.
Parâmetro
Descrição
Obrigatório
Observações
metastore
O tipo de metastore.
Sim
Neste exemplo, selecione DLF.
catalog name
O nome do catálogo de dados do DLF.
ImportanteSe você usar um usuário RAM ou uma função RAM, certifique-se de que o usuário RAM ou a função RAM tenha as permissões necessárias para ler e gravar dados no DLF. Para obter mais informações, consulte Gerenciar permissões de dados.
Sim
Recomendamos que você use o DLF 2.5. Esta versão elimina a necessidade de inserir informações como um par AccessKey e permite selecionar rapidamente um catálogo de dados DLF existente. Para obter informações sobre como criar um catálogo de dados, consulte Catálogo de dados.
Após criar um catálogo de dados chamado paimoncatalog, selecione-o na lista.
-
No catálogo de dados, crie um banco de dados chamado order_dw para sincronizar todas as tabelas do banco de dados order_dw do MySQL.
No painel de navegação à esquerda, escolha e clique em New para abrir uma janela de consulta temporária.
-- Use the paimoncatalog data source USE CATALOG paimoncatalog; -- Create the order_dw database CREATE DATABASE order_dw;A mensagem retornada
The following statement has been executed successfully!indica que o banco de dados foi criado com sucesso.
Para obter mais informações sobre como usar catálogos Paimon, consulte Gerenciar catálogos Paimon.
Criar um catálogo MySQL
Na página Metadata Management, clique em Create Catalog.
Na aba Built-in Catalog, clique em MySQL e, em seguida, clique em Next.
-
Configure os parâmetros a seguir e clique em OK para criar um catálogo MySQL chamado mysqlcatalog.
Parâmetro
Descrição
Obrigatório
Observações
catalog name
O nome do catálogo.
Sim
Insira um nome personalizado. Este tutorial usa mysqlcatalog.
hostname
O endereço IP ou nome do host do banco de dados MySQL.
Sim
Para obter mais informações, consulte Visualizar e gerenciar endpoints e portas de conexão da instância. Como a instância ApsaraDB RDS for MySQL e o workspace do Realtime Compute for Apache Flink estão na mesma VPC, insira o endpoint interno.
port
O número da porta do serviço de banco de dados MySQL. O valor padrão é 3306.
Não
Para obter mais informações, consulte Visualizar e gerenciar endpoints e portas de conexão da instância.
default-database
O nome do banco de dados MySQL padrão.
Sim
Insira o nome do banco de dados a ser sincronizado, que é order_dw neste tutorial.
username
O nome de usuário para o serviço de banco de dados MySQL.
Sim
Use a conta criada em Preparar a fonte de dados MySQL CDC.
password
A senha para o serviço de banco de dados MySQL.
Sim
Use a senha da conta criada em Preparar a fonte de dados MySQL CDC.
Construir a camada ODS: Ingestão em tempo real
Construa a camada ODS em uma única etapa usando um job de ingestão de dados Flink CDC no formato yaml para sincronizar dados do MySQL para o Paimon.
-
Crie e inicie o job de ingestão de dados yaml.
No console do Realtime Compute for Apache Flink, acesse a página e crie um job rascunho yaml em branco chamado ods.
-
Copie o código a seguir para o editor e atualize parâmetros como nome de usuário e senha.
source: type: mysql name: MySQL Source hostname: rm-bp1e********566g.mysql.rds.aliyuncs.com port: 3306 username: ${secret_values.username} password: ${secret_values.password} tables: order_dw.\.* # Use a regular expression to read all tables in the order_dw database. # (Optional) Synchronize data from newly created tables during the incremental phase. scan.binlog.newly-added-table.enabled: true # (Optional) Synchronize table and field comments. include-comments.enabled: true # (Optional) Prioritize the distribution of unbounded chunks to prevent potential TaskManager OutOfMemory issues. scan.incremental.snapshot.unbounded-chunk-first.enabled: true # (Optional) Enable parsing filters to accelerate reads. scan.only.deserialize.captured.tables.changelog.enabled: true sink: type: paimon name: Paimon Sink catalog.properties.metastore: rest catalog.properties.uri: http://cn-beijing-vpc.dlf.aliyuncs.com catalog.properties.warehouse: paimoncatalog catalog.properties.token.provider: dlf pipeline: name: MySQL to Paimon PipelineParâmetro
Descrição
Obrigatório
Exemplo
catalog.properties.metastoreO tipo de metastore. Defina como rest.
Sim
rest
catalog.properties.token.providerO provedor de token. Defina como DLF.
Sim
DLF
catalog.properties.uriA URI de acesso para o DLF Rest Catalog Server é
http://[region-id]-vpc.dlf.aliyuncs.com. Para obter mais informações sobre o ID da Região, consulte Endpoints de Serviço.Sim
http://cn-beijing-vpc.dlf.aliyuncs.com
catalog.properties.warehouseO nome do catálogo DLF.
Sim
paimoncatalog
hostnameO endereço IP ou nome do host do banco de dados MySQL. Para obter mais informações, consulte Visualizar e gerenciar endpoints e portas de conexão da instância. Como a instância ApsaraDB RDS for MySQL e o workspace do Realtime Compute for Apache Flink estão na mesma VPC, insira o endpoint interno.
Sim
rm-bp1e****566g.mysql.rds.aliyuncs.com
usernameO nome de usuário do banco de dados MySQL. Recomendamos usar o gerenciamento de segredos. Para obter mais informações, consulte Gerenciar variáveis.
Sim
${secret_values.username}
passwordA senha do banco de dados MySQL. Recomendamos usar o gerenciamento de segredos. Para obter mais informações, consulte Gerenciar variáveis.
Sim
${secret_values.password}
Para otimizar o desempenho de gravação do Paimon, consulte Ajuste de desempenho do Paimon.
No canto superior direito, clique em Deploy.
Em , clique em Start na coluna Actions do job ODS recém-implantado e selecione stateless start para iniciar o job. Para obter mais informações sobre configurações de inicialização de jobs, consulte Inicialização de Job.
-
Visualize os dados nas três tabelas que foram sincronizadas do MySQL para o Paimon.
No console do Realtime Compute for Apache Flink, acesse a página . Na aba Query Editor, copie o código a seguir para o editor, selecione a instrução e clique em Run no canto superior direito.
SELECT * FROM paimoncatalog.order_dw.orders ORDER BY order_id;Após a execução da consulta, ela retorna sete registros de pedidos da tabela orders. O resultado inclui as colunas order_id (100001 a 100007), user_id, shop_id, product_id, buy_fee (1000 a 5000), create_time, update_time e state. O valor para state é 1 para todos os registros.
Construir a camada DWD: tabela larga temática
-
Crie a tabela larga dwd_orders
No console do Realtime Compute for Apache Flink, acesse a página . Na aba Query Editor, copie o código a seguir para o editor, selecione a instrução e clique em Run no canto superior direito.
CREATE TABLE paimoncatalog.order_dw.dwd_orders ( order_id BIGINT, order_user_id STRING, order_shop_id BIGINT, order_product_id BIGINT, order_product_catalog_name STRING, order_fee BIGINT, order_create_time TIMESTAMP, order_update_time TIMESTAMP, order_state INT, pay_id BIGINT, pay_platform INT COMMENT 'platform 0: phone, 1: pc', pay_create_time TIMESTAMP, PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'merge-engine' = 'partial-update', -- Use the partial-update merge engine to generate the wide table 'changelog-producer' = 'lookup' -- Use the lookup changelog producer to generate changelogs with low latency );A mensagem
The following statement has been executed successfully!indica uma criação bem-sucedida. -
Consuma dados de alteração do ODS
No console do Realtime Compute for Apache Flink, acesse a página . Crie um novo job de streaming SQL chamado dwd, copie o código a seguir para o editor SQL, clique em Deploy e inicie o job usando Stateless Start.
Este job SQL une a tabela orders com a tabela product_catalog em uma junção de tabela de dimensão. Em seguida, ele grava o resultado na tabela dwd_orders junto com dados da tabela orders_pay. O mecanismo de mesclagem de atualização parcial do Paimon une registros das tabelas orders e orders_pay que possuem o mesmo order_id.
SET 'execution.checkpointing.max-concurrent-checkpoints' = '3'; SET 'table.exec.sink.upsert-materialize' = 'NONE'; SET 'execution.checkpointing.interval' = '10s'; SET 'execution.checkpointing.min-pause' = '10s'; -- Paimon does not currently support multiple INSERT statements into the same table within a single job. Therefore, UNION ALL is used here. INSERT INTO paimoncatalog.order_dw.dwd_orders SELECT o.order_id, o.user_id, o.shop_id, o.product_id, dim.catalog_name, o.buy_fee, o.create_time, o.update_time, o.state, NULL, NULL, NULL FROM paimoncatalog.order_dw.orders o LEFT JOIN paimoncatalog.order_dw.product_catalog FOR SYSTEM_TIME AS OF proctime() AS dim ON o.product_id = dim.product_id UNION ALL SELECT order_id, NULL, NULL, NULL, NULL, NULL, NULL, NULL, NULL, pay_id, pay_platform, create_time FROM paimoncatalog.order_dw.orders_pay; -
Visualize os dados em dwd_orders
No console do Realtime Compute for Apache Flink, acesse a página . Na aba Query Editor, copie o código a seguir para o editor, selecione a instrução e clique em Run no canto superior direito.
SELECT * FROM paimoncatalog.order_dw.dwd_orders ORDER BY order_id;Após a execução bem-sucedida da consulta, ela retorna sete registros de pedidos da tabela larga
dwd_orders. O resultado inclui oito campos:order_id,order_user_id,order_shop_id,order_product_id,order_product_catalog_name,order_fee,order_create_timeeorder_update_time. Oorder_idvaria de 100001 a 100007, as taxas de pedido variam de 1000 a 5000 e todos os horários de criação são em 2023-02-15.
Construir a camada DWS: cálculo de métricas
-
Crie as tabelas agregadas da camada DWS dws_users e dws_shops.
No console do Realtime Compute for Apache Flink, acesse a página . Na aba Query Editor, copie o código a seguir para o editor, selecione a instrução e clique em Run no canto superior direito.
-- User dimension aggregate table. CREATE TABLE paimoncatalog.order_dw.dws_users ( user_id STRING, ds STRING, paid_buy_fee_sum BIGINT COMMENT 'Total amount paid on the current day', PRIMARY KEY (user_id, ds) NOT ENFORCED ) WITH ( 'merge-engine' = 'aggregation', -- Use the aggregation merge engine to generate the aggregate table. 'fields.paid_buy_fee_sum.aggregate-function' = 'sum' -- Aggregate the result by summing the data in paid_buy_fee_sum. -- Because the dws_users table is not consumed by downstream streaming jobs, you do not need to specify a changelog producer. ); -- Shop dimension aggregate table. CREATE TABLE paimoncatalog.order_dw.dws_shops ( shop_id BIGINT, ds STRING, paid_buy_fee_sum BIGINT COMMENT 'Total amount paid on the current day', uv BIGINT COMMENT 'Total number of distinct purchasing users on the current day', pv BIGINT COMMENT 'Total number of purchases by users on the current day', PRIMARY KEY (shop_id, ds) NOT ENFORCED ) WITH ( 'merge-engine' = 'aggregation', -- Use the aggregation merge engine to generate the aggregate table. 'fields.paid_buy_fee_sum.aggregate-function' = 'sum', -- Aggregate the result by summing the data in paid_buy_fee_sum. 'fields.uv.aggregate-function' = 'sum', -- Aggregate the result by summing the data in uv. 'fields.pv.aggregate-function' = 'sum' -- Aggregate the result by summing the data in pv. -- Because the dws_shops table is not consumed by downstream streaming jobs, you do not need to specify a changelog producer. ); -- To calculate both user- and shop-dimension aggregate tables simultaneously, create an intermediate table with a primary key of user and shop. CREATE TABLE paimoncatalog.order_dw.dwm_users_shops ( user_id STRING, shop_id BIGINT, ds STRING, paid_buy_fee_sum BIGINT COMMENT 'Total amount a user paid in a shop on the current day', pv BIGINT COMMENT 'Number of times a user purchased from a shop on the current day', PRIMARY KEY (user_id, shop_id, ds) NOT ENFORCED ) WITH ( 'merge-engine' = 'aggregation', -- Use the aggregation merge engine to generate the aggregate table. 'fields.paid_buy_fee_sum.aggregate-function' = 'sum', -- Aggregate the result by summing the data in paid_buy_fee_sum. 'fields.pv.aggregate-function' = 'sum', -- Aggregate the result by summing the data in pv. 'changelog-producer' = 'lookup', -- Use the lookup changelog producer to generate changelogs with low latency. -- Intermediate tables in the DWM layer are not typically queried directly by upstream applications, so they can be optimized for write performance. 'file.format' = 'avro', -- Using the Avro row-based storage format provides more efficient write performance. 'metadata.stats-mode' = 'none' -- Omitting statistics increases the cost of OLAP queries (with no impact on continuous stream processing) but improves write performance. );A mensagem
The following statement has been executed successfully!indica que a criação foi bem-sucedida. -
Consuma dados de alteração da tabela dwd_orders da camada DWD.
No console do Realtime Compute for Apache Flink, acesse a página , crie um novo job de streaming SQL chamado dwm. Copie o código a seguir para o editor SQL, clique em Deploy e inicie o job usando Stateless Start.
Este job SQL grava dados da tabela dwd_orders na tabela dwm_users_shops. O mecanismo de mesclagem de dados de pré-agregação do Paimon soma automaticamente o order_fee para calcular o valor total gasto por um usuário em uma loja. Ele também soma o valor 1 para calcular o número de vezes que um usuário fez uma compra nessa loja.
SET 'execution.checkpointing.max-concurrent-checkpoints' = '3'; SET 'table.exec.sink.upsert-materialize' = 'NONE'; SET 'execution.checkpointing.interval' = '10s'; SET 'execution.checkpointing.min-pause' = '10s'; INSERT INTO paimoncatalog.order_dw.dwm_users_shops SELECT order_user_id, order_shop_id, DATE_FORMAT (pay_create_time, 'yyyyMMdd') as ds, order_fee, 1 -- One input record represents one purchase FROM paimoncatalog.order_dw.dwd_orders WHERE pay_id IS NOT NULL AND order_fee IS NOT NULL; -
Consuma dados de alteração da tabela dwm_users_shops da camada DWM em tempo real.
No console do Realtime Compute for Apache Flink, acesse a página . Crie um novo job de streaming SQL chamado dws. Copie o código a seguir para o editor SQL, clique em Deploy e inicie o job usando Stateless Start.
Este job SQL grava dados da tabela dwm_users_shops nas tabelas dws_users e dws_shops. O mecanismo de mesclagem de dados de pré-agregação do Paimon é usado para calcular o gasto total de cada usuário (paid_buy_fee_sum) na tabela dws_users. Na tabela dws_shops, ele calcula as vendas totais da loja (paid_buy_fee_sum), o número de usuários compradores (somando 1) e o número total de compras (pv).
SET 'execution.checkpointing.max-concurrent-checkpoints' = '3'; SET 'table.exec.sink.upsert-materialize' = 'NONE'; SET 'execution.checkpointing.interval' = '10s'; SET 'execution.checkpointing.min-pause' = '10s'; -- Unlike the DWD job, each INSERT statement here writes to a different Paimon table, so they can be included in the same job. BEGIN STATEMENT SET; INSERT INTO paimoncatalog.order_dw.dws_users SELECT user_id, ds, paid_buy_fee_sum FROM paimoncatalog.order_dw.dwm_users_shops; -- The primary key is the shop ID. Data for popular shops might be much larger than for others. -- Therefore, use local merge to pre-aggregate data in memory before writing to Paimon to mitigate data skew. INSERT INTO paimoncatalog.order_dw.dws_shops /*+ OPTIONS('local-merge-buffer-size' = '64mb') */ SELECT shop_id, ds, paid_buy_fee_sum, 1, -- One input record represents all of a user's purchases at this shop pv FROM paimoncatalog.order_dw.dwm_users_shops; END; -
Visualize os dados nas tabelas dws_users e dws_shops.
No console do Realtime Compute for Apache Flink, acesse a página . Na aba Query Editor, copie o código a seguir para o editor, selecione a instrução e clique em Run no canto superior direito.
--View the data in the dws_users table SELECT * FROM paimoncatalog.order_dw.dws_users ORDER BY user_id;Após a execução da consulta, a tabela dws_users retorna três registros com as seguintes colunas:
user_id,dsepaid_buy_fee_sum. Os registros são:user_001/ 20230215 / 8000,user_002/ 20230215 / 5000 euser_003/ 20230215 / 5000.--View the data in the dws_shops table SELECT * FROM paimoncatalog.order_dw.dws_shops ORDER BY shop_id;Ao consultar a tabela dws_shops, o resultado contém quatro linhas e cinco colunas:
shop_id,ds,paid_buy_fee_sum,uvepv. Os shop_ids são 12345 a 12348, ds é 20230215 para todos os registros, os valores de paid_buy_fee_sum são 5000, 4000, 7000 e 2000, os valores de uv são 1, 1, 2 e 2, e os valores de pv são 1, 1, 3 e 2, respectivamente.
Capturar alterações do banco de dados de negócios
Após construir o streaming data lakehouse, teste sua capacidade de capturar alterações do banco de dados de negócios.
-
Insira os seguintes dados no banco de dados order_dw no MySQL.
INSERT INTO orders VALUES (100008, 'user_001', 12345, 3, 3000, '2023-02-15 17:40:56', '2023-02-15 18:42:56', 1), (100009, 'user_002', 12348, 4, 1000, '2023-02-15 18:40:56', '2023-02-15 19:42:56', 1), (100010, 'user_003', 12348, 2, 2000, '2023-02-15 19:40:56', '2023-02-15 20:42:56', 1); INSERT INTO orders_pay VALUES (2008, 100008, 1, '2023-02-15 18:40:56'), (2009, 100009, 1, '2023-02-15 19:40:56'), (2010, 100010, 0, '2023-02-15 20:40:56'); -
Visualize os dados nas tabelas dws_users e dws_shops. No console do Realtime Compute for Apache Flink, acesse a página . Na aba Query Editor, copie o código a seguir para o editor, selecione a instrução e clique em Run no canto superior direito.
-
Tabela dws_users
SELECT * FROM paimoncatalog.order_dw.dws_users ORDER BY user_id;A consulta retorna três linhas com as seguintes colunas:
user_id,dsepaid_buy_fee_sum. Os registros são:user_001/ 20230215 / 11000,user_002/ 20230215 / 6000 euser_003/ 20230215 / 7000. -
Tabela dws_shops
SELECT * FROM paimoncatalog.order_dw.dws_shops ORDER BY shop_id;A consulta retorna quatro linhas com cinco colunas:
shop_id,ds,paid_buy_fee_sum,uvepv. Os shop_ids são 12345, 12346, 12347 e 12348, ds é 20230215 para todos os registros, os valores de paid_buy_fee_sum são 8000, 4000, 7000 e 5000, os valores de uv são 1, 1, 2 e 3, e os valores de pv são 2, 1, 3 e 4, respectivamente.
-
Usar o streaming data lakehouse
A seção anterior demonstrou como criar um catálogo Paimon e gravar em tabelas Paimon no Flink. Esta seção mostra alguns casos de uso simples de análise de dados usando StarRocks após a configuração do streaming data lakehouse.
Conectar StarRocks ao DLF
Para obter mais informações, consulte Acessar o DLF a partir do Serverless StarRocks.
Consulta de classificação
A seguinte consulta do StarRocks analisa a tabela agregada na camada DWS para recuperar as três principais lojas por valor de transação em 15 de fevereiro de 2023.
SELECT ROW_NUMBER() OVER (ORDER BY paid_buy_fee_sum DESC) AS rn, shop_id, paid_buy_fee_sum
FROM paimoncatalog.order_dw.dws_shops
WHERE ds = '20230215'
ORDER BY rn LIMIT 3;
A consulta retorna três linhas: rn=1 corresponde a shop_id=12345 e paid_buy_fee_sum=8000.00; rn=2 corresponde a shop_id=12347 e paid_buy_fee_sum=7000.00; e rn=3 corresponde a shop_id=12348 e paid_buy_fee_sum=5000.00.
Consulta de detalhes
A seguinte consulta do StarRocks analisa a tabela larga na camada DWD para recuperar os detalhes do pedido de um cliente específico que usou uma plataforma de pagamento específica em fevereiro de 2023.
SELECT * FROM paimoncatalog.order_dw.dwd_orders
WHERE order_create_time >= '2023-02-01 00:00:00' AND order_create_time < '2023-03-01 00:00:00'
AND order_user_id = 'user_001'
AND pay_platform = 0
ORDER BY order_create_time;;
A consulta retorna o seguinte resultado.
order_id order_user_id order_shop_id order_product_id order_product_catalog_name order_fee
0 100006 user_001 12348 1 phone_aaa 1000
1 100004 user_001 12347 4 phone_ddd 2000
Relatório de dados
A seguinte consulta do StarRocks analisa a tabela larga na camada DWD para gerar um relatório sobre o número total de pedidos e o valor total dos pedidos para cada data e categoria de produto em fevereiro de 2023.
SELECT
order_create_time AS order_create_date,
order_product_catalog_name,
COUNT(*),
SUM(order_fee)
FROM
paimoncatalog.order_dw.dwd_orders
WHERE
order_create_time >= '2023-02-01 00:00:00' and order_create_time < '2023-03-01 00:00:00'
GROUP BY
order_create_date, order_product_catalog_name
ORDER BY
order_create_date, order_product_catalog_name;
A consulta retorna nove registros, mostrando dados de pedidos por hora das 10:40:56 às 18:40:56 em 15 de fevereiro de 2023. As categorias de produtos incluem phone_ddd, phone_aaa, phone_eee, phone_ccc e phone_bbb. Cada registro tem um count(*) de 1, e os valores de sum(order_fee) são 2000, 1000, 1000, 2000, 3000, 4000, 5000, 3000 e 1000, respectivamente.
Referências
Para construir um lakehouse offline Paimon usando processamento em lote do Flink, consulte Processamento em lote do Flink.