ApsaraDB for SelectDB é totalmente compatível com o Apache Doris. Use o Flink Doris Connector para importar dados históricos de fontes como MySQL, Oracle, PostgreSQL, SQL Server e Kafka para o SelectDB. Após iniciar uma tarefa de captura de dados de alteração (CDC) no Flink, ele também sincroniza dados incrementais da fonte para o SelectDB.
Visão geral
Atualmente, o Flink Doris Connector suporta apenas a gravação de dados no SelectDB. Caso precise usar o Flink Doris Connector para se conectar diretamente aos nós de back-end do SelectDB e ler dados com eficiência, entre em contato com a equipe de suporte técnico do SelectDB para solicitar acesso.
Você também pode usar o Flink JDBC Connector para ler dados do SelectDB.
O Flink Doris Connector permite que o Flink leia e escreva no Apache Doris para processamento e análise de dados em tempo real. Como o SelectDB possui total compatibilidade com o Apache Doris, este connector é um método comum para ingestão de dados em streaming no SelectDB.
Cada componente funciona da seguinte maneira:
-
source
Finalidade: Uma source lê dados de sistemas externos para um fluxo de dados do Flink. Esses sistemas podem incluir filas de mensagens (como Apache Kafka), bancos de dados ou sistemas de arquivos.
Exemplo: Use o Kafka como source para ler mensagens em tempo real ou leia dados de um arquivo.
-
Transform
Finalidade: O estágio de transformação processa o fluxo de dados recebido. Essas operações podem incluir filtragem, mapeamento, agregação e janelas (windowing).
Exemplo: Mapeie um fluxo de entrada para converter sua estrutura de dados ou agregue dados para calcular uma métrica por minuto.
-
Sink
Finalidade: Um sink grava os dados processados de um fluxo do Flink em um sistema externo, como banco de dados, arquivo ou fila de mensagens.
Exemplo: Grave os resultados processados em um banco de dados MySQL ou envie os dados para outro tópico do Kafka.
A figura a seguir ilustra como os dados são importados para o SelectDB usando o Flink Doris Connector.
Pré-requisitos
-
Garanta a conectividade de rede entre sua fonte de dados, o Flink e o SelectDB.
-
Solicite um endpoint público para sua instância do ApsaraDB for SelectDB. Para mais informações, consulte Solicitar ou liberar um endpoint público.
Pule esta etapa se o ambiente Flink e a fonte de dados estiverem na mesma Virtual Private Cloud (VPC) da sua instância do ApsaraDB for SelectDB. Isso é comum quando são produtos da Alibaba Cloud ou estão implantados em instâncias do Elastic Compute Service (ECS) na mesma VPC.
Adicione os endereços IP do seu ambiente Flink e da fonte de dados à lista de permissões da sua instância do ApsaraDB for SelectDB. Para mais informações, consulte Configurar uma lista de permissões de endereços IP.
-
-
Certifique-se de que o Flink Doris Connector esteja instalado.
A tabela a seguir lista os requisitos de versão para o Flink e o Flink Doris Connector.
Versão do Flink
Versão do Flink Doris Connector
Link para download
Realtime Compute for Apache Flink: 1,17 ou posterior
Flink open source: 1,15 ou posterior
1.5.2 ou posterior. Recomendamos baixar a versão mais recente.
Para instruções de instalação, consulte Instalar o Flink Doris Connector.
Adicionar o Flink Doris Connector
Adicione o Flink Doris Connector de acordo com o seu ambiente.
Se você usa o Realtime Compute for Apache Flink para importar dados no SelectDB, gerencie o
Flink Doris Connectorcomo um connector personalizado. Para detalhes, consulte Gerenciar connectors personalizados.Caso utilize um cluster Flink autogerenciado, baixe o pacote JAR correspondente do
Flink Doris Connectore coloque-o no diretóriolibda sua instalação do Flink. Para o link de download, consulte Pacote JAR.-
Para adicionar o
Flink Doris Connectorcomo dependência Maven, insira o código abaixo no arquivo de configuração de dependências do seu projeto. Para outras versões, consulte o Repositório Maven.<!-- flink-doris-connector --> <dependency> <groupId>org.apache.doris</groupId> <artifactId>flink-doris-connector-1.16</artifactId> <version>1.5.2</version> </dependency>
Exemplos
Ambiente de exemplo
Este exemplo utiliza Flink SQL, Flink CDC e a API DataStream para migrar dados da tabela employees no banco de dados test de uma instância ApsaraDB RDS for MySQL para a tabela employees no banco de dados test de uma instância do SelectDB. Modifique os parâmetros nestes exemplos para adequá-los ao seu cenário. O ambiente de exemplo é o seguinte:
Ambiente standalone do Flink 1.16
Java
Banco de dados de destino: test
Tabela de destino: employees
Banco de dados de origem: test
Tabela de origem: employees
Preparar o ambiente
Ambiente Flink
-
Prepare um ambiente Java.
O Flink requer um ambiente Java para execução. Instale um Java Development Kit (JDK) e configure a variável de ambiente
JAVA_HOME.Para obter uma lista das versões Java suportadas, consulte Compatibilidade com Java. Este exemplo usa o Java 8. Para instruções de instalação, consulte Instalar o JDK.
-
Baixe o pacote de instalação do Flink flink-1.16.3-bin-scala_2.12.tgz. Se esta versão estiver desatualizada, baixe outra versão em Apache Flink.
wget https://www.apache.si/flink/flink-1.16.3/flink-1.16.3-bin-scala_2.12.tgz -
Descompacte o pacote de instalação.
tar -zxvf flink-1.16.3-bin-scala_2.12.tgz -
Acesse o diretório
libno diretório de instalação do Flink e adicione os connectors necessários para as etapas seguintes.-
Adicione o Flink Doris Connector.
wget https://repo.maven.apache.org/maven2/org/apache/doris/flink-doris-connector-1.16/1.5.2/flink-doris-connector-1.16-1.5.2.jar -
Adicione o Flink MySQL Connector.
wget https://repo1.maven.org/maven2/com/ververica/flink-sql-connector-mysql-cdc/2.4.2/flink-sql-connector-mysql-cdc-2.4.2.jar
-
-
Inicie o cluster Flink.
No diretório
binda sua instalação do Flink, execute o seguinte comando:./start-cluster.sh
SelectDB de destino
Crie uma instância do ApsaraDB for SelectDB. Para mais informações, consulte Criar uma instância.
Conecte-se à instância. Para mais informações, consulte Conectar-se a uma instância.
-
Crie um banco de dados de teste chamado
test.CREATE DATABASE test; -
Crie uma tabela de teste chamada
employees.USE test; -- Create table CREATE TABLE employees ( emp_no int NOT NULL, birth_date date, first_name varchar(20), last_name varchar(20), gender char(2), hire_date date ) UNIQUE KEY(`emp_no`) DISTRIBUTED BY HASH(`emp_no`) BUCKETS 1;
MySQL de origem
Crie uma instância do ApsaraDB RDS for MySQL.
-
Crie um banco de dados de teste chamado
test.CREATE DATABASE test; -
Crie uma tabela de teste chamada
employees.USE test; CREATE TABLE employees ( emp_no INT NOT NULL PRIMARY KEY, birth_date DATE, first_name VARCHAR(20), last_name VARCHAR(20), gender CHAR(2), hire_date DATE ); -
Insira dados.
INSERT INTO employees (emp_no, birth_date, first_name, last_name, gender, hire_date) VALUES (1001, '1985-05-15', 'John', 'Doe', 'M', '2010-06-20'), (1002, '1990-08-22', 'Jane', 'Smith', 'F', '2012-03-15'), (1003, '1987-11-02', 'Robert', 'Johnson', 'M', '2015-07-30'), (1004, '1992-01-18', 'Emily', 'Davis', 'F', '2018-01-05'), (1005, '1980-12-09', 'Michael', 'Brown', 'M', '2008-11-21');
Importar com Flink SQL
-
Inicie o Flink SQL Client.
No diretório
binda sua instalação do Flink, execute o seguinte comando:./sql-client.sh -
No Flink SQL Client, envie um job do Flink.
-
Crie uma tabela de origem MySQL.
A cláusula
WITHna instrução a seguir especifica a configuração para aMySQL CDC Source. Para mais informações sobre os parâmetros, consulte MySQL | Apache Flink CDC.CREATE TABLE employees_source ( emp_no INT, birth_date DATE, first_name STRING, last_name STRING, gender STRING, hire_date DATE, PRIMARY KEY (`emp_no`) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = '127.0.0.1', 'port' = '3306', 'username' = 'root', 'password' = '****', 'database-name' = 'test', 'table-name' = 'employees' ); -
Crie uma tabela sink do SelectDB.
A cláusula
WITHna instrução a seguir especifica a configuração para o SelectDB. Para mais informações sobre os parâmetros, consulte Parâmetros do Sink.CREATE TABLE employees_sink ( emp_no INT , birth_date DATE, first_name STRING, last_name STRING, gender STRING, hire_date DATE ) WITH ( 'connector' = 'doris', 'fenodes' = 'selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080', 'table.identifier' = 'test.employees', 'username' = 'admin', 'password' = '****' ); -
Sincronize dados da tabela de origem MySQL para a tabela sink do SelectDB.
INSERT INTO employees_sink SELECT * FROM employees_source;
-
-
Verifique a importação de dados.
Conecte-se ao SelectDB e execute a seguinte instrução para visualizar os dados importados.
SELECT * FROM test.employees;
Importar com Flink CDC
O Realtime Compute for Apache Flink não suporta jobs baseados em JAR. Em vez disso, utilize jobs baseados em YAML com o CDC 3.0.
Use o Flink CDC para importar dados no SelectDB.
Para executar um job do Flink CDC, utilize o programa flink no diretório de instalação do Flink. A sintaxe é a seguinte:
<FLINK_HOME>/bin/flink run \
-Dexecution.checkpointing.interval=10s \
-Dparallelism.default=1 \
-c org.apache.doris.flink.tools.cdc.CdcTools \
lib/flink-doris-connector-1.16-1.5.2.jar \
<mysql-sync-database|oracle-sync-database|postgres-sync-database|sqlserver-sync-database> \
--database <selectdb-database-name> \
[--job-name <flink-job-name>] \
[--table-prefix <selectdb-table-prefix>] \
[--table-suffix <selectdb-table-suffix>] \
[--including-tables <mysql-table-name|name-regular-expr>] \
[--excluding-tables <mysql-table-name|name-regular-expr>] \
--mysql-conf <mysql-cdc-source-conf> [--mysql-conf <mysql-cdc-source-conf> ...] \
--oracle-conf <oracle-cdc-source-conf> [--oracle-conf <oracle-cdc-source-conf> ...] \
--sink-conf <doris-sink-conf> [--table-conf <doris-sink-conf> ...] \
[--table-conf <selectdb-table-conf> [--table-conf <selectdb-table-conf> ...]]
Parâmetros
|
Parâmetro |
Descrição |
||
|
execution.checkpointing.interval |
Intervalo de checkpoint do Flink. Esta configuração afeta a frequência de sincronização de dados. Recomenda-se o valor de |
||
|
parallelism.default |
Paralelismo do job do Flink. Aumentar o paralelismo pode melhorar a velocidade de sincronização de dados. |
||
|
job-name |
Nome do job do Flink. |
||
|
database |
Nome do banco de dados de destino no SelectDB. |
||
|
table-prefix |
Prefixo para o nome da tabela de destino no SelectDB. Por exemplo, |
||
|
table-suffix |
Sufixo para o nome da tabela de destino no SelectDB. |
||
|
including-tables |
Tabelas a serem sincronizadas. Use uma barra vertical |
` para separar várias tabelas. Expressões regulares são suportadas. Por exemplo, |
tbl.*` sincroniza |
|
excluding-tables |
Tabelas a serem excluídas da sincronização. O formato é o mesmo de |
||
|
mysql-conf |
Configuração para a MySQL CDC Source. Para mais informações, consulte MySQL CDC Connector. Os parâmetros |
||
|
oracle-conf |
Configuração para a Oracle CDC Source. Para mais informações, consulte Oracle CDC Connector. Os parâmetros |
||
|
sink-conf |
Parâmetros de configuração para o Doris Sink. Para mais informações, consulte Parâmetros do Sink. |
||
|
table-conf |
Parâmetros de configuração para a tabela do SelectDB. Estes correspondem ao conteúdo da cláusula |
Para sincronização de dados, adicione a dependência necessária do Flink CDC, como flink-sql-connector-mysql-cdc-${version}.jar ou flink-sql-connector-oracle-cdc-${version}.jar, ao diretório $FLINK_HOME/lib.
A sincronização completa de banco de dados é suportada no Flink 1.15 e versões posteriores. Para baixar diferentes versões do Flink Doris Connector, consulte Flink Doris Connector.
Parâmetros do Sink
|
Parâmetro |
Padrão |
Obrigatório |
Descrição |
|
fenodes |
Nenhum |
Sim |
O endpoint e a porta HTTP da sua instância do ApsaraDB for SelectDB. Você pode obter o VPC Endpoint (ou Public Endpoint) e a HTTP Port na página Instance Details > Network Information no console do ApsaraDB for SelectDB. Exemplo: |
|
table.identifier |
Nenhum |
Sim |
Nome do banco de dados e da tabela. Exemplo: |
|
username |
Nenhum |
Sim |
Nome de usuário do banco de dados da sua instância do ApsaraDB for SelectDB. |
|
password |
Nenhum |
Sim |
Senha do usuário do banco de dados da sua instância do ApsaraDB for SelectDB. |
|
jdbc-url |
Nenhum |
Não |
Informações de conexão JDBC da sua instância do ApsaraDB for SelectDB. Você pode obter o VPC Endpoint (ou Public Endpoint) e a MySQL Port na página Instance Details > Network Information no console do ApsaraDB for SelectDB. Exemplo: |
|
auto-redirect |
true |
Não |
Define se deve redirecionar requisições Stream Load. Se ativado, o Stream Load grava dados através dos frontends (FEs) e as informações do backend (BE) não são recuperadas. |
|
doris.request.retries |
3 |
Não |
Número de tentativas de reenvio de requisição para o SelectDB. |
|
doris.request.connect.timeout |
30s |
Não |
Tempo limite para conexão com o SelectDB. |
|
doris.request.read.timeout |
30s |
Não |
Tempo limite para leitura de dados do SelectDB. |
|
sink.label-prefix |
"" |
Sim |
Prefixo do rótulo para importações via Stream Load. Em cenários de two-phase commit (2PC), este prefixo deve ser globalmente único para garantir a semântica exactly-once (EOS) do Flink. |
|
sink.properties |
Nenhum |
Não |
Parâmetros de importação para o Stream Load. Configure as propriedades da seguinte forma:
Para mais parâmetros, consulte Stream Load. |
|
sink.buffer-size |
1048576 |
Não |
Tamanho do buffer de gravação, em bytes. Recomenda-se o valor padrão de 1 MB. |
|
sink.buffer-count |
3 |
Não |
Número de buffers de gravação. Recomenda-se o valor padrão. |
|
sink.max-retries |
3 |
Não |
Número máximo de novas tentativas após falha no commit. O padrão é 3. |
|
sink.use-cache |
false |
Não |
Define se deve usar cache em memória para recuperação em caso de exceções. Se ativado, os dados do período de checkpoint são mantidos no cache. |
|
sink.enable-delete |
true |
Não |
Define se deve sincronizar eventos de exclusão. Esta opção é suportada apenas para tabelas que usam o modelo Unique Key. |
|
sink.enable-2pc |
true |
Não |
Define se deve ativar o two-phase commit (2PC). Ativado por padrão ( |
|
sink.enable.batch-mode |
false |
Não |
Define se deve usar o modo batch para gravar dados no SelectDB. Quando ativado, a operação de gravação é acionada pelo tamanho do buffer ou tempo, conforme definido por Quando o modo batch está ativado, a semântica exactly-once (EOS) não é garantida. Você pode usar o modelo Unique Key para alcançar idempotência. |
|
sink.flush.queue-size |
2 |
Não |
Tamanho da fila de buffer no modo batch. |
|
sink.buffer-flush.max-rows |
50000 |
Não |
Número máximo de linhas por gravação em lote no modo batch. |
|
sink.buffer-flush.max-bytes |
10MB |
Não |
Tamanho máximo em bytes por gravação em lote no modo batch. |
|
sink.buffer-flush.interval |
10s |
Não |
Intervalo assíncrono de limpeza do buffer no modo batch. O valor mínimo é 1 segundo. |
|
sink.ignore.update-before |
true |
Não |
Define se deve ignorar eventos |
Exemplos de sincronização
Sincronização MySQL
<FLINK_HOME>/bin/flink run \
-Dexecution.checkpointing.interval=10s \
-Dparallelism.default=1 \
-c org.apache.doris.flink.tools.cdc.CdcTools \
lib/flink-doris-connector-1.16-1.5.2.jar \
mysql-sync-database \
--database test \
--mysql-conf hostname=127.0.0.1 \
--mysql-conf port=3306 \
--mysql-conf username=root \
--mysql-conf password="password" \
--mysql-conf database-name=test \
--including-tables "employees" \
--sink-conf fenodes=selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080 \
--sink-conf username=admin \
--sink-conf password=****
Sincronização Oracle
<FLINK_HOME>/bin/flink run \
-Dexecution.checkpointing.interval=10s \
-Dparallelism.default=1 \
-c org.apache.doris.flink.tools.cdc.CdcTools \
lib/flink-doris-connector-1.16-1.5.2.jar \
oracle-sync-database \
--database test_db \
--oracle-conf hostname=127.0.0.1 \
--oracle-conf port=1521 \
--oracle-conf username=admin \
--oracle-conf password="password" \
--oracle-conf database-name=XE \
--oracle-conf schema-name=ADMIN \
--including-tables "tbl1|test.*" \
--sink-conf fenodes=selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080 \
--sink-conf username=admin \
--sink-conf password=****
Sincronização PostgreSQL
<FLINK_HOME>/bin/flink run \
-Dexecution.checkpointing.interval=10s \
-Dparallelism.default=1 \
-c org.apache.doris.flink.tools.cdc.CdcTools \
lib/flink-doris-connector-1.16-1.5.2.jar \
postgres-sync-database \
--database db1\
--postgres-conf hostname=127.0.0.1 \
--postgres-conf port=5432 \
--postgres-conf username=postgres \
--postgres-conf password="123456" \
--postgres-conf database-name=postgres \
--postgres-conf schema-name=public \
--postgres-conf slot.name=test \
--postgres-conf decoding.plugin.name=pgoutput \
--including-tables "tbl1|test.*" \
--sink-conf fenodes=selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080 \
--sink-conf username=admin \
--sink-conf password=****
Sincronização SQL Server
<FLINK_HOME>/bin/flink run \
-Dexecution.checkpointing.interval=10s \
-Dparallelism.default=1 \
-c org.apache.doris.flink.tools.cdc.CdcTools \
lib/flink-doris-connector-1.16-1.5.2.jar \
sqlserver-sync-database \
--database db1\
--sqlserver-conf hostname=127.0.0.1 \
--sqlserver-conf port=1433 \
--sqlserver-conf username=sa \
--sqlserver-conf password="123456" \
--sqlserver-conf database-name=CDC_DB \
--sqlserver-conf schema-name=dbo \
--including-tables "tbl1|test.*" \
--sink-conf fenodes=selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080 \
--sink-conf username=admin \
--sink-conf password=****
Importar com a API DataStream
-
Adicione as seguintes dependências ao seu projeto Maven.
Dependências Maven
-
Código Java principal.
O código a seguir configura a tabela de origem MySQL e a tabela sink do ApsaraDB for SelectDB. Os parâmetros correspondem aos usados na seção Importar dados usando Flink SQL. Para mais informações, consulte MySQL | Apache Flink CDC e Parâmetros do Sink.
package org.example; import com.ververica.cdc.connectors.mysql.source.MySqlSource; import com.ververica.cdc.connectors.mysql.table.StartupOptions; import com.ververica.cdc.connectors.shaded.org.apache.kafka.connect.json.JsonConverterConfig; import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema; import org.apache.doris.flink.cfg.DorisExecutionOptions; import org.apache.doris.flink.cfg.DorisOptions; import org.apache.doris.flink.sink.DorisSink; import org.apache.doris.flink.sink.writer.serializer.JsonDebeziumSchemaSerializer; import org.apache.doris.flink.tools.cdc.mysql.DateToStringConverter; import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.streaming.api.datastream.DataStreamSource; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import java.util.HashMap; import java.util.Map; import java.util.Properties; public class Main { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); env.enableCheckpointing(10000); Map<String, Object> customConverterConfigs = new HashMap<>(); customConverterConfigs.put(JsonConverterConfig.DECIMAL_FORMAT_CONFIG, "numeric"); JsonDebeziumDeserializationSchema schema = new JsonDebeziumDeserializationSchema(false, customConverterConfigs); // Configure the MySQL source table MySqlSource<String> mySqlSource = MySqlSource.<String>builder() .hostname("rm-xxx.mysql.rds.aliyuncs***") .port(3306) .startupOptions(StartupOptions.initial()) .databaseList("db_test") .tableList("db_test.employees") .username("root") .password("test_123") .debeziumProperties(DateToStringConverter.DEFAULT_PROPS) .deserializer(schema) .serverTimeZone("Asia/Shanghai") .build(); // Configure the ApsaraDB for SelectDB sink table DorisSink.Builder<String> sinkBuilder = DorisSink.builder(); DorisOptions.Builder dorisBuilder = DorisOptions.builder(); dorisBuilder.setFenodes("selectdb-cn-xxx-public.selectdbfe.rds.aliyunc****:8080") .setTableIdentifier("db_test.employees") .setUsername("admin") .setPassword("test_123"); DorisOptions dorisOptions = dorisBuilder.build(); // Configure Stream Load parameters with sink.properties Properties properties = new Properties(); properties.setProperty("format", "json"); properties.setProperty("read_json_by_line", "true"); DorisExecutionOptions.Builder executionBuilder = DorisExecutionOptions.builder(); executionBuilder.setStreamLoadProp(properties); sinkBuilder.setDorisExecutionOptions(executionBuilder.build()) .setSerializer(JsonDebeziumSchemaSerializer.builder().setDorisOptions(dorisOptions).build()) // Serialize the data stream. .setDorisOptions(dorisOptions); DataStreamSource<String> dataStreamSource = env.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source"); dataStreamSource.sinkTo(sinkBuilder.build()); env.execute("MySQL to SelectDB"); } }
Uso avançado
Atualizar colunas parciais com Flink SQL
-- enable checkpoint
SET 'execution.checkpointing.interval' = '10s';
CREATE TABLE cdc_mysql_source (
id INT
,name STRING
,bank STRING
,age INT
,PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = '127.0.0.1',
'port' = '3306',
'username' = 'root',
'password' = 'password',
'database-name' = 'database',
'table-name' = 'table'
);
CREATE TABLE selectdb_sink (
id INT,
name STRING,
bank STRING,
age INT
)
WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'database.table',
'username' = 'admin',
'password' = '****',
'sink.properties.format' = 'json',
'sink.properties.read_json_by_line' = 'true',
'sink.properties.columns' = 'id,name,bank,age',
'sink.properties.partial_columns' = 'true' -- Enable partial column updates.
);
INSERT INTO selectdb_sink SELECT id,name,bank,age FROM cdc_mysql_source;
Usar Flink SQL para excluir dados por coluna
Em cenários de CDC, o sink do Doris identifica o tipo de evento a partir de RowKind e atribui um valor à coluna oculta __DORIS_DELETE_SIGN__ para realizar exclusões. Quando a fonte de dados são mensagens do Kafka, o sink não consegue usar RowKind para determinar o tipo de operação. Em vez disso, ele deve confiar em um campo específico dentro da mensagem, como {"op_type":"delete",data:{...}}. Para excluir dados onde op_type é 'delete', você deve passar explicitamente um valor para a coluna oculta com base na sua lógica de negócios. O exemplo de Flink SQL a seguir mostra como excluir dados no Alibaba Cloud SelectDB com base em um campo específico nos dados do Kafka.
-- Example message: {"op_type":"delete",data:{"id":1,"name":"zhangsan"}}
CREATE TABLE KAFKA_SOURCE(
data STRING,
op_type STRING
) WITH (
'connector' = 'kafka',
...
);
CREATE TABLE SELECTDB_SINK(
id INT,
name STRING,
__DORIS_DELETE_SIGN__ INT
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'db.table',
'username' = 'admin',
'password' = '****',
'sink.enable-delete' = 'false', -- A value of false indicates that the event type is not inferred from RowKind.
'sink.properties.columns' = 'id, name, __DORIS_DELETE_SIGN__' -- Explicitly specify the columns for the Stream Load import.
);
INSERT INTO SELECTDB_SINK
SELECT json_value(data,'$.id') as id,
json_value(data,'$.name') as name,
if(op_type='delete',1,0) as __DORIS_DELETE_SIGN__
FROM KAFKA_SOURCE;
FAQ
-
P: Como escrevo dados
BITMAP?R: Veja o exemplo abaixo:
CREATE TABLE bitmap_sink ( dt INT, page STRING, user_id INT ) WITH ( 'connector' = 'doris', 'fenodes' = 'selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080', 'table.identifier' = 'test.bitmap_test', 'username' = 'admin', 'password' = '****', 'sink.label-prefix' = 'selectdb_label', 'sink.properties.columns' = 'dt,page,user_id,user_id=to_bitmap(user_id)' ); -
P: Como resolvo o erro
errCode = 2, detailMessage = Label[label_0_1]has already been used, relate to txn[19650]?R: Em um cenário exactly-once, um job do Flink deve ser reiniciado a partir do checkpoint ou savepoint mais recente. Este erro ocorre se você reiniciar o job a partir de um estado mais antigo. Se a semântica exactly-once não for necessária, desative o two-phase commit (2PC) definindo
sink.enable-2pc=falseou use um sink.label-prefix diferente. -
P: Como resolvo o erro
errCode = 2, detailMessage = transaction[19650]not found?R: Este erro ocorre durante a fase de commit. Indica que o ID da transação registrado no checkpoint expirou no ApsaraDB for SelectDB. Quando o connector tenta confirmar essa transação expirada, o servidor relata que a transação não foi encontrada. Nesse caso, não é possível reiniciar o job a partir do checkpoint. Para evitar esse problema, aumente o parâmetro
streaming_label_keep_max_secondno ApsaraDB for SelectDB. O valor padrão é 12 horas. -
P: Como resolvo o erro
errCode = 2, detailMessage = current running txns on db 10006 is 100, larger than limit 100?R: Este erro indica que o número de transações de importação simultâneas para um único banco de dados excedeu o limite do sistema de 100. Para resolver, aumente o parâmetro
max_running_txn_num_per_dbno ApsaraDB for SelectDB. Para mais informações, consulte max_running_txn_num_per_db.Esse erro também pode ocorrer se você alterar frequentemente o rótulo e reiniciar o job. Em cenários de two-phase commit (2PC) (aplicáveis aos modelos Duplicate Key e Aggregate Key), cada job requer um rótulo único. Quando um job reinicia a partir de um checkpoint, o Flink aborta apenas as transações que foram pré-confirmadas (precommitted), mas ainda não confirmadas (committed). Se você alterar frequentemente o rótulo antes de reiniciar, muitas transações pré-confirmadas não são abortadas e continuam consumindo a cota de transações. Para o modelo Unique Key, desative o 2PC e projete o operador sink para gravações idempotentes.
-
P: Como garanto a ordenação de dados dentro de um lote ao gravar em uma tabela que usa o modelo Unique Key?
R: Adicione uma configuração de coluna de sequência para garantir a ordenação dos dados. Para mais informações, consulte SEQUENCE.
-
P: Por que nenhum dado está sendo sincronizado mesmo que o job do Flink não relate erros?
R: Esse comportamento depende da versão do connector. Em versões anteriores à 1.1.0, as gravações são em lote e orientadas por dados, portanto, verifique se a fonte upstream está produzindo dados. Na versão 1.1.0 e posteriores, as gravações são acionadas por checkpoints, que devem estar habilitados para gravar dados.
-
P: Como resolvo o erro
tablet writer write failed, tablet_id=190958, txn_id=3505530, err=-235?R: Este erro geralmente ocorre em versões do connector anteriores à 1.1.0. É causado por uma frequência de gravação excessivamente alta, que cria muitas versões no tablet. Para resolver, aumente os parâmetros
sink.buffer-flush.max-bytesesink.buffer-flush.intervalpara reduzir a frequência do Stream Load. -
P: Como ignoro dados incorretos (dirty data) durante uma importação Flink?
R: Se os dados de origem contiverem registros que não correspondem ao esquema da tabela de destino (por exemplo, tipo de dados ou comprimento incorreto), o job Stream Load falha e o Flink tenta novamente continuamente. Para ignorar esses dados incorretos, desative o modo estrito do Stream Load definindo
strict_mode=false,max_filter_ratio=1, ou adicione uma etapa de transformação para filtrar os dados inválidos antes que cheguem ao operador sink. -
P: Como a tabela de origem deve ser mapeada para a tabela do ApsaraDB for SelectDB?
R: Ao importar dados usando o Flink Doris Connector, garanta que os seguintes mapeamentos estejam corretos: (1) As colunas e tipos na tabela de origem devem corresponder aos do Flink SQL. (2) As colunas e tipos no Flink SQL devem corresponder aos da tabela do ApsaraDB for SelectDB.
-
P: Como resolvo o erro
TApplicationException: get_next failed: out of sequence response: expected 4 but got 3?R: Este erro indica um bug de concorrência no framework Thrift subjacente. Para resolver, atualize para a versão mais recente do Flink Doris Connector e use uma versão compatível do Flink.
-
P: Como resolvo o erro
DorisRuntimeException: Fail to abort transaction 26153 with urlhttp://192.168.XX.XX?R: Para diagnosticar este problema, pesquise nos logs do TaskManager pela frase
abort transaction response. O código de status HTTP na entrada do log indica se o problema se origina do cliente ou do servidor.