A integração do MaxCompute com o Kafka oferece recursos eficientes e confiáveis de processamento e análise de dados. Essa integração é ideal para cenários que exigem processamento em tempo real, fluxos de dados em grande escala e análises complexas. Este tópico descreve como gravar dados do Message Queue for Apache Kafka e de instâncias Kafka autogerenciadas no MaxCompute, além de fornecer um exemplo detalhado para uma instância Kafka autogerenciada.
Gravar dados do Kafka no MaxCompute: Kafka totalmente gerenciado pela Alibaba Cloud
O MaxCompute possui integração nativa com o Message Queue for Apache Kafka. Utilize o MaxCompute Sink Connector for Message Queue for Apache Kafka para importar continuamente dados de um tópico específico para uma tabela do MaxCompute, sem necessidade de ferramentas de terceiros ou desenvolvimento personalizado. Para mais informações, consulte Crie um conector sink do MaxCompute.
Gravar dados do Kafka no MaxCompute: Kafka open source autogerenciado
Pré-requisitos
Implante o Kafka V2.2 ou posterior e crie um tópico do Kafka. Recomenda-se a versão 3.4.0.
Crie um projeto e uma tabela no MaxCompute. Para mais informações, consulte Crie um projeto do MaxCompute e Crie uma tabela.
Observações
O conector do Kafka suporta a gravação de dados nos formatos TEXT, CSV, JSON e FLATTEN. As observações abaixo se aplicam a cada formato. Para mais informações sobre tipos de dados, consulte Descrições dos tipos de dados.
-
Ao gravar dados do Kafka nos formatos TEXT ou JSON no MaxCompute, a tabela do MaxCompute deve atender aos seguintes requisitos:
Nome do campo
Tipo do campo
Campo fixo
topic
STRING
Sim
partition
BIGINT
Sim
offset
BIGINT
Sim
key
Para gravação de dados Kafka em formato TEXT, o tipo do campo deve ser STRING.
Para gravação de dados Kafka em formato JSON, o tipo do campo pode ser STRING ou JSON, dependendo dos dados a serem gravados.
Este campo é fixo para sincronizar a chave da mensagem do Kafka com a tabela do MaxCompute. Para mais informações sobre o modo de sincronização de mensagens do Kafka para o MaxCompute, consulte modo.
value
Para gravação de dados Kafka em formato TEXT, o tipo do campo deve ser STRING.
Para gravação de dados Kafka em formato JSON, o tipo do campo pode ser STRING ou JSON, dependendo dos dados a serem gravados.
Este campo é fixo para sincronizar o valor da mensagem do Kafka com a tabela do MaxCompute. Para mais informações sobre o modo de sincronização de mensagens do Kafka para o MaxCompute, consulte modo.
pt
STRING (campo de partição)
Sim
-
Para gravar dados do Kafka nos formatos FLATTEN ou CSV no MaxCompute, a tabela deve incluir os campos e tipos de dados listados abaixo. Defina outros campos conforme os dados que você está gravando.
Nome do campo
Tipo do campo
topic
STRING
partition
BIGINT
offset
BIGINT
pt
STRING (campo de partição)
Na gravação de dados Kafka em formato CSV para uma tabela do MaxCompute, a ordem e os tipos de dados dos campos personalizados na tabela devem corresponder às colunas nos dados do Kafka para garantir o sucesso da operação.
-
Na gravação de dados Kafka em formato FLATTEN para uma tabela do MaxCompute, os nomes dos campos personalizados na tabela devem corresponder aos nomes dos campos nos dados do Kafka para garantir o sucesso da operação.
Por exemplo, se os dados Kafka em formato FLATTEN forem
{"A":a,"B":"b","C":{"D":"d","E":"e"}}, configure a tabela do MaxCompute da seguinte forma.CREATE TABLE IF NOT EXISTS table_flatten( topic STRING, `partition` BIGINT, `offset` BIGINT, A BIGINT, B STRING, C JSON ) PARTITIONED BY (pt STRING);
Configure e inicie o serviço do conector Kafka
-
Este exemplo utiliza um ambiente Linux. Em uma janela de comando, baixe o pacote
kafka-connector-2.0.jarexecutando o comando abaixo ou usando o link de download.wget http://maxcompute-repo.oss-cn-hangzhou.aliyuncs.com/kafka/kafka-connector-2.0.jarPara evitar conflitos de dependência, crie uma subpasta, como
connector, no diretório$KAFKA_HOME/libse coloque o pacotekafka-connector-2.0.jardentro dela.NotaCaso o pacote
kafka-connector-2.0.jarnão seja compatível com seu ambiente de implantação do Kafka, consulte Configure Kafka-connector para obter mais informações sobre como configurar e iniciar o serviçoKafka-connector. -
No diretório
$KAFKA_HOME/config, configure o arquivoconnect-distributed.properties.Adicione o conteúdo a seguir ao arquivo
connect-distributed.properties.## Add the following content plugin.path=<KAFKA_HOME>/libs/connector ## Update the values of the key.converter and value.converter parameters key.converter=org.apache.kafka.connect.storage.StringConverter value.converter=org.apache.kafka.connect.storage.StringConverter -
No diretório
$KAFKA_HOME/, execute o comando abaixo para iniciar o serviçoKafka-connector.## Start command bin/connect-distributed.sh config/connect-distributed.properties &
Configure e inicie a tarefa do conector Kafka
-
Crie e configure o arquivo de configuração
odps-sink-connector.json. Em seguida, carregue o arquivoodps-sink-connector.jsonem qualquer local.As seções a seguir descrevem o conteúdo e os parâmetros do arquivo de configuração
odps-sink-connector.json.{ "name": "Kafka connector task name", "config": { "connector.class": "com.aliyun.odps.kafka.connect.MaxComputeSinkConnector", "tasks.max": "3", "topics": "your_topic", "endpoint": "endpoint", "tunnel_endpoint": "your_tunnel endpoint", "project": "project", "schema":"default", "table": "your_table", "account_type": "account type (STS or ALIYUN)", "access_id": "access id", "access_key": "access key", "account_id": "account id for sts", "sts.endpoint": "sts endpoint", "region_id": "region id for sts", "role_name": "role name for sts", "client_timeout_ms": "STS Token valid period (ms)", "format": "TEXT", "mode": "KEY", "partition_window_type": "MINUTE", "use_streaming": false, "buffer_size_kb": 65536, "sink_pool_size":"150", "record_batch_size":"8000", "runtime.error.topic.name":"kafka topic when runtime errors happens", "runtime.error.topic.bootstrap.servers":"kafka bootstrap servers of error topic queue", "skip_error":"false" } }-
Parâmetros comuns
Parâmetro
Obrigatório
Descrição
name
Sim
Nome da tarefa. O nome deve ser único.
connector.class
Sim
Nome da classe para iniciar o serviço
Kafka connector. O valor padrão écom.aliyun.odps.kafka.connect.MaxComputeSinkConnector.tasks.max
Sim
Número máximo de processos consumidores no
Kafka connector. O valor deve ser um número inteiro maior que 0.topics
Sim
Nome do tópico do Kafka.
endpoint
Sim
Endpoint do serviço MaxCompute.
Configure o endpoint com base na região e no tipo de conectividade de rede selecionados durante a criação do projeto MaxCompute. Para consultar os endpoints de cada região e rede, veja Endpoints.
tunnel_endpoint
Não
Endpoint público do serviço Tunnel.
Se você não configurar um endpoint do Tunnel, o roteamento será feito automaticamente para o endpoint correspondente à rede onde o serviço MaxCompute está localizado. Caso configure um endpoint, sua definição terá precedência e o roteamento automático será desativado.
Para consultar os endpoints de cada região e rede, veja Endpoints.
project
Sim
Nome do projeto MaxCompute de destino.
schema
Não
Parâmetro obrigatório se o projeto MaxCompute de destino estiver configurado com modelo de schema de três camadas. O valor padrão é default.
Parâmetro não obrigatório se o projeto MaxCompute de destino não estiver configurado com modelo de schema de três camadas.
Para mais informações sobre schemas, consulte Operações de schema.
table
Sim
Nome da tabela no projeto MaxCompute de destino.
format
Não
Formato das mensagens a serem gravadas. Valores válidos:
TEXT (padrão): A mensagem é uma string.
BINARY: A mensagem é um array de bytes.
CSV: A mensagem é uma string com valores separados por vírgulas (,).
JSON: A mensagem é uma string no tipo de dados JSON. Para mais informações sobre o tipo JSON do MaxCompute, consulte Tipo de dados JSON.
FLATTEN: A mensagem é uma string no tipo de dados JSON. As chaves e valores na string JSON são analisados e gravados nas colunas correspondentes da tabela do MaxCompute. As chaves nos dados JSON devem corresponder aos nomes das colunas na tabela do MaxCompute.
Para exemplos de importação de mensagens em diferentes formatos, consulte Exemplos de uso.
mode
Não
Modo de sincronização de mensagens para o MaxCompute. Valores válidos:
KEY: Mantém apenas a chave da mensagem e a grava na tabela MaxCompute de destino.
VALUE: Mantém apenas o valor da mensagem e o grava na tabela MaxCompute de destino.
DEFAULT (padrão): Mantém tanto a chave quanto o valor da mensagem e os grava na tabela MaxCompute de destino.
No modo DEFAULT, apenas os formatos de dados TEXT e BINARY são suportados.
partition_window_type
Não
Particiona dados por hora do sistema. Valores válidos: DAY, HOUR (padrão) e MINUTE.
use_streaming
Não
Define se o túnel de dados streaming será usado. Valores válidos:
false (padrão): Desativado.
true: Ativado.
buffer_size_kb
Não
Tamanho do buffer interno do gravador de partições odps, em KB. O valor padrão é 65536 KB.
sink_pool_size
Não
Número máximo de threads para gravação multithread. O valor padrão corresponde ao número de núcleos de CPU do sistema.
record_batch_size
Não
Número máximo de mensagens que uma única thread dentro de uma tarefa do conector Kafka pode enviar em paralelo de uma só vez.
skip_error
Não
Define se registros que causam erros desconhecidos serão ignorados. Valores válidos:
false (padrão): Não ignora os registros.
true: Ignora os registros.
NotaSe skip_error estiver definido como false e o parâmetro runtime.error.topic.name não estiver configurado, o processo interrompe a gravação de dados quando ocorre um erro desconhecido. O processo fica bloqueado e uma exceção é lançada no log.
Se skip_error estiver definido como true e runtime.error.topic.name não estiver configurado, o processo de gravação continua e os dados anômalos são descartados.
Se skip_error estiver definido como false e runtime.error.topic.name estiver configurado, o processo de gravação continua e os dados anômalos são registrados no tópico especificado por runtime.error.topic.name.
Para um exemplo de como lidar com dados anômalos, consulte Exemplo de tratamento de dados anômalos.
runtime.error.topic.name
Não
Nome do tópico do Kafka onde serão gravados os dados que causarem erros desconhecidos durante a operação de gravação.
runtime.error.topic.bootstrap.servers
Não
Endereço do servidor bootstrap da instância Kafka onde serão gravados os dados que causarem erros desconhecidos durante a operação de gravação.
account_type
Sim
Método de acesso ao serviço MaxCompute de destino. Os valores válidos são STS e ALIYUN. O valor padrão é ALIYUN.
Diferentes métodos de acesso exigem parâmetros de credenciais distintos. Para mais informações, consulte Acessar o MaxCompute usando o método ALIYUN e Acessar o MaxCompute usando o método STS.
-
Além dos parâmetros comuns, configure também os parâmetros a seguir.
Nome do Parâmetro
Descrição
access_id
AccessKey ID da sua conta Alibaba Cloud ou usuário RAM.
Obtenha o AccessKey ID na página AccessKey Management.
access_key
AccessKey secret correspondente ao AccessKey ID.
-
Além dos parâmetros comuns, configure também os parâmetros a seguir.
Parâmetro
Descrição
account_id
ID da conta usada para acessar o projeto MaxCompute de destino. Visualize o ID da sua conta no Account Center.
region_id
ID da região do projeto MaxCompute de destino. Para consultar o ID de cada região, veja Endpoints.
role_name
Nome da função usada para acessar o projeto MaxCompute de destino. Visualize o nome da função na página Roles.
client_timeout_ms
Intervalo de atualização do token do Security Token Service (STS), em milissegundos (ms). O valor padrão é 11 ms.
sts.endpoint
Endpoint do serviço STS necessário para autenticação de identidade usando um token de segurança temporário (STS).
Para consultar os endpoints de cada região e rede, veja Endpoints.
-
-
Execute o comando abaixo para iniciar a tarefa de migração de dados do conector Kafka.
curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" http://localhost:8083/connectors -d @odps-sink-connector.json
Exemplos de uso
Gravar dados TEXT
-
Prepare os dados.
-
Use um cliente local (odpscmd) ou outra ferramenta capaz de executar comandos SQL do MaxCompute para se conectar ao MaxCompute e criar a tabela de destino.
CREATE TABLE IF NOT EXISTS table_text( topic STRING, `partition` BIGINT, `offset` BIGINT, key STRING, value STRING ) PARTITIONED BY (pt STRING); -
Crie os dados do Kafka.
No diretório
$KAFKA_HOME/bin/, execute o comando abaixo para criar um tópico do Kafka. Este exemplo usatopic_textcomo nome do tópico.sh kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic topic_textExecute o comando abaixo para criar mensagens do Kafka.
sh kafka-console-producer.sh --bootstrap-server localhost:9092 --topic topic_text --property parse.key=true >123 abc >456 edf
-
-
(Opcional) Inicie o serviço
Kafka-connector. Para mais informações, consulte Configure e inicie o serviço do conector Kafka.NotaSe o serviço
Kafka-connectorjá estiver em execução, pule esta etapa. -
Crie e configure o arquivo
odps-sink-connector.json. Em seguida, carregue o arquivoodps-sink-connector.jsonem qualquer local, como o caminho$KAFKA_HOME/config.O código abaixo fornece um exemplo do arquivo
odps-sink-connector.json. Para mais informações sobre o arquivoodps-sink-connector.json, consulte Configure e inicie a tarefa do conector Kafka.{ "name": "odps-test-text", "config": { "connector.class": "com.aliyun.odps.kafka.connect.MaxComputeSinkConnector", "tasks.max": "3", "topics": "topic_text", "endpoint": "http://service.cn-shanghai.maxcompute.aliyun.com/api", "project": "project_name", "schema":"default", "table": "table_text", "account_type": "ALIYUN", "access_id": "<yourAccessKeyId>", "access_key": "<yourAccessKeySecret>", "partition_window_type": "MINUTE", "mode":"VALUE", "format":"TEXT", "sink_pool_size":"150", "record_batch_size":"9000", "buffer_size_kb":"600000" } } -
Execute o comando abaixo para iniciar a tarefa de migração de dados do conector Kafka.
curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" http://localhost:8083/connectors -d @$KAFKA_HOME/config/odps-sink-connector.json -
Verifique o resultado.
Use um cliente local (odpscmd) ou outra ferramenta capaz de executar comandos SQL do MaxCompute para se conectar ao MaxCompute e, em seguida, execute o comando abaixo para consultar os dados e verificar o resultado.
set odps.sql.allow.fullscan=true; select * from table_text;A saída retornada é:
# Because the mode parameter in the odps-sink-connector.json configuration file is set to VALUE, only the content of the value is retained. The key field is NULL. +-------+------------+------------+-----+-------+----+ | topic | partition | offset | key | value | pt | +-------+------------+------------+-----+-------+----+ | topic_text | 0 | 0 | NULL | abc | 07-13-2023 21:13 | | topic_text | 0 | 1 | NULL | edf | 07-13-2023 21:13 | +-------+------------+------------+-----+-------+----+
Gravar dados CSV
-
Prepare os dados.
-
Use um cliente local (odpscmd) ou outra ferramenta capaz de executar comandos SQL do MaxCompute para se conectar ao MaxCompute e criar a tabela de destino.
CREATE TABLE IF NOT EXISTS table_csv( topic STRING, `partition` BIGINT, `offset` BIGINT, id BIGINT, name STRING, region STRING ) PARTITIONED BY (pt STRING); -
Grave dados no Kafka.
No diretório
$KAFKA_HOME/bin/, execute o comando abaixo para criar um tópico do Kafka chamadotopic_csv.sh kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic topic_csvExecute o comando abaixo para criar mensagens do Kafka.
sh kafka-console-producer.sh --bootstrap-server localhost:9092 --topic topic_csv --property parse.key=true >123 1103,zhangsan,china >456 1104,lisi,usa
-
-
(Opcional) Inicie o serviço
Kafka-connector. Para mais informações, consulte Configure e inicie o serviço do conector Kafka.NotaSe o serviço
Kafka-connectorjá estiver em execução, pule esta etapa. -
Crie e configure o arquivo
odps-sink-connector.jsone, em seguida, carregue o arquivoodps-sink-connector.jsonem qualquer local. Este tópico usa o caminho$KAFKA_HOME/configcomo exemplo.O código abaixo fornece um exemplo do arquivo
odps-sink-connector.json. Para mais informações sobre o arquivoodps-sink-connector.json, consulte Configure e inicie a tarefa do conector Kafka.{ "name": "odps-test-csv", "config": { "connector.class": "com.aliyun.odps.kafka.connect.MaxComputeSinkConnector", "tasks.max": "3", "topics": "topic_csv", "endpoint": "http://service.cn-shanghai.maxcompute.aliyun.com/api", "project": "project_name", "schema":"default", "table": "table_csv", "account_type": "ALIYUN", "access_id": "<yourAccessKeyId>", "access_key": "<yourAccessKeySecret>", "partition_window_type": "MINUTE", "format":"CSV", "mode":"VALUE", "sink_pool_size":"150", "record_batch_size":"9000", "buffer_size_kb":"600000" } } -
Execute o comando abaixo para iniciar a tarefa de migração de dados do conector Kafka.
curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" http://localhost:8083/connectors -d @$KAFKA_HOME/config/odps-sink-connector.json -
Verifique o resultado.
Use um cliente local (odpscmd) ou outra ferramenta capaz de executar comandos SQL do MaxCompute para se conectar ao MaxCompute e, em seguida, execute o comando abaixo para consultar os dados e verificar o resultado.
set odps.sql.allow.fullscan=true; select * from table_csv;A saída retornada é:
+-------+------------+------------+------------+------+--------+----+ | topic | partition | offset | id | name | region | pt | +-------+------------+------------+------------+------+--------+----+ | csv_test | 0 | 0 | 1103 | zhangsan | china | 07-14-2023 00:10 | | csv_test | 0 | 1 | 1104 | lisi | usa | 07-14-2023 00:10 | +-------+------------+------------+------------+------+--------+----+
Gravar dados JSON
-
Prepare os dados.
-
Use um cliente local (odpscmd) ou outra ferramenta capaz de executar comandos SQL do MaxCompute para se conectar ao MaxCompute e criar a tabela de destino.
CREATE TABLE IF NOT EXISTS table_json( topic STRING, `partition` BIGINT, `offset` BIGINT, key STRING, value JSON ) PARTITIONED BY (pt STRING); -
Crie os dados do Kafka.
No diretório
$KAFKA_HOME/bin/, execute o comando abaixo para criar um tópico do Kafka. Este exemplo usatopic_jsoncomo nome do tópico.sh kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic topic_jsonExecute o comando abaixo para criar mensagens do Kafka.
sh kafka-console-producer.sh --bootstrap-server localhost:9092 --topic topic_json --property parse.key=true >123 {"id":123,"name":"json-1","region":"beijing"} >456 {"id":456,"name":"json-2","region":"hangzhou"}
-
-
(Opcional) Inicie o serviço
Kafka-connector. Para mais informações, consulte Configure e inicie o serviço do conector Kafka.NotaSe o serviço
Kafka-connectorjá estiver em execução, pule esta etapa. -
Crie e configure o arquivo
odps-sink-connector.json. Em seguida, carregue o arquivoodps-sink-connector.jsonem qualquer local, como o caminho$KAFKA_HOME/config.O código abaixo fornece um exemplo do arquivo
odps-sink-connector.json. Para mais informações sobre o arquivoodps-sink-connector.json, consulte Configure e inicie a tarefa do conector Kafka.{ "name": "odps-test-json", "config": { "connector.class": "com.aliyun.odps.kafka.connect.MaxComputeSinkConnector", "tasks.max": "3", "topics": "topic_json", "endpoint": "http://service.cn-shanghai.maxcompute.aliyun.com/api", "project": "project_name", "schema":"default", "table": "table_json", "account_type": "ALIYUN", "access_id": "<yourAccessKeyId>", "access_key": "<yourAccessKeySecret>", "partition_window_type": "MINUTE", "mode":"VALUE", "format":"JSON", "sink_pool_size":"150", "record_batch_size":"9000", "buffer_size_kb":"600000" } } -
Execute o comando abaixo para iniciar a tarefa de migração de dados do conector Kafka.
curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" http://localhost:8083/connectors -d @$KAFKA_HOME/config/odps-sink-connector.json -
Verifique o resultado.
Use um cliente local (odpscmd) ou outra ferramenta capaz de executar comandos SQL do MaxCompute para se conectar ao MaxCompute e, em seguida, execute o comando abaixo para consultar os dados e verificar o resultado.
set odps.sql.allow.fullscan=true; select * from table_json;A saída retornada é:
# The JSON data is successfully written to the value field. +-------+------------+------------+-----+-------+----+ | topic | partition | offset | key | value | pt | +-------+------------+------------+-----+-------+----+ | Topic_json | 0 | 0 | NULL | {"id":123,"name":"json-1","region":"beijing"} | 07-14-2023 00:28 | | Topic_json | 0 | 1 | NULL | {"id":456,"name":"json-2","region":"hangzhou"} | 07-14-2023 00:28 | +-------+------------+------------+-----+-------+----+
Gravar dados FLATTEN
-
Prepare os dados.
-
Use um cliente local (odpscmd) ou outra ferramenta capaz de executar comandos SQL do MaxCompute para se conectar ao MaxCompute e criar a tabela de destino.
CREATE TABLE IF NOT EXISTS table_flatten( topic STRING, `partition` BIGINT, `offset` BIGINT, id BIGINT, name STRING, extendinfo JSON ) PARTITIONED BY (pt STRING); -
Crie os dados do Kafka.
No diretório
$KAFKA_HOME/bin/, execute o comando abaixo para criar um tópico do Kafka. Este exemplo usatopic_flattencomo nome do tópico../kafka/bin/kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic topic_flattenExecute o comando abaixo para criar mensagens do Kafka.
sh kafka-console-producer.sh --bootstrap-server localhost:9092 --topic topic_flatten --property parse.key=true >123 {"id":123,"name":"json-1","extendinfo":{"region":"beijing","sex":"M"}} >456 {"id":456,"name":"json-2","extendinfo":{"region":"hangzhou","sex":"W"}}
-
-
(Opcional) Inicie o serviço
Kafka-connector. Para mais informações, consulte Configure e inicie o serviço do conector Kafka.NotaSe o serviço
Kafka-connectorjá estiver em execução, pule esta etapa. -
Crie e configure o arquivo
odps-sink-connector.jsone, em seguida, carregue o arquivoodps-sink-connector.jsonem qualquer local. Este tópico usa o caminho$KAFKA_HOME/configcomo exemplo.O código abaixo fornece um exemplo do arquivo
odps-sink-connector.json. Para mais informações sobre o arquivoodps-sink-connector.json, consulte Configure e inicie a tarefa do conector Kafka.{ "name": "odps-test-flatten", "config": { "connector.class": "com.aliyun.odps.kafka.connect.MaxComputeSinkConnector", "tasks.max": "3", "topics": "topic_flatten", "endpoint": "http://service.cn-shanghai.maxcompute.aliyun.com/api", "project": "project_name", "schema":"default", "table": "table_flatten", "account_type": "ALIYUN", "access_id": "<yourAccessKeyId>", "access_key": "<yourAccessKeySecret>", "partition_window_type": "MINUTE", "mode":"VALUE", "format":"FLATTEN", "sink_pool_size":"150", "record_batch_size":"9000", "buffer_size_kb":"600000" } } -
Execute o comando abaixo para iniciar a tarefa do conector Kafka.
curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" http://localhost:8083/connectors -d @$KAFKA_HOME/config/odps-sink-connector.json -
Verifique o resultado.
Use um cliente local (odpscmd) ou outra ferramenta capaz de executar comandos SQL do MaxCompute para se conectar ao MaxCompute e, em seguida, execute o comando abaixo para consultar os dados e verificar o resultado.
set odps.sql.allow.fullscan=true; select * from table_flatten;O resultado é mostrado abaixo:
# The JSON data is parsed and written to a MaxCompute table, with extendinfo as a JSON field that supports nesting. +-------+------------+--------+-----+------+------------+----+ | topic | partition | offset | id | name | extendinfo | pt | +-------+------------+--------+-----+------+------------+----+ | topic_flatten | 0 | 0 | 123 | json-1 | {"sex":"M","region":"beijing"} | 07-14-2023 01:33 | | topic_flatten | 0 | 1 | 456 | json-2 | {"sex":"W","region":"hangzhou"} | 07-14-2023 01:33 | +-------+------------+--------+-----+------+------------+----+
Exemplo de tratamento de dados anômalos
-
Prepare os dados.
-
Use um cliente local (odpscmd) ou outra ferramenta capaz de executar comandos SQL do MaxCompute para se conectar ao MaxCompute e criar a tabela de destino.
CREATE TABLE IF NOT EXISTS table_flatten( topic STRING, `partition` BIGINT, `offset` BIGINT, id BIGINT, name STRING, extendinfo JSON ) PARTITIONED BY (pt STRING); -
Crie os dados do Kafka.
No diretório
$KAFKA_HOME/bin/, execute os comandos abaixo para criar tópicos do Kafka.-
O tópico
topic_abnormal.sh kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic topic_abnormal -
Tópico de mensagens para exceções
runtime_error.sh kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic runtime_errorNotaSe ocorrer um erro durante uma operação de gravação de dados, os dados anômalos serão gravados no tópico
runtime_error. Esse tipo de erro geralmente é causado por incompatibilidade entre os dados do Kafka e o schema da tabela do MaxCompute.
Execute o comando abaixo para criar mensagens do Kafka.
Uma das mensagens no comando abaixo não corresponde ao schema da tabela MaxCompute de destino.
sh kafka-console-producer.sh --bootstrap-server localhost:9092 --topic flatten_test --property parse.key=true >100 {"id":100,"name":"json-3","extendinfo":{"region":"beijing","gender":"M"}} >101 {"id":101,"name":"json-4","extendinfos":"null"} >102 {"id":102,"name":"json-5","extendinfo":{"region":"beijing","gender":"M"}} -
-
-
(Opcional) Inicie o serviço
Kafka-connector. Para mais informações, consulte Configure e inicie o serviço do conector Kafka.NotaSe o serviço
Kafka-connectorjá estiver em execução, pule esta etapa. -
Crie e configure o arquivo
odps-sink-connector.jsone, em seguida, carregue o arquivoodps-sink-connector.jsonem qualquer local. Este tópico usa o caminho$KAFKA_HOME/configcomo exemplo.O código abaixo fornece um exemplo do arquivo
odps-sink-connector.json. Para mais informações sobre o arquivoodps-sink-connector.json, consulte Configure e inicie a tarefa do conector Kafka.{ "name": "odps-test-runtime-error", "config": { "connector.class": "com.aliyun.odps.kafka.connect.MaxComputeSinkConnector", "tasks.max": "3", "topics": "topic_abnormal", "endpoint": "http://service.cn-shanghai.maxcompute.aliyun.com/api", "project": "project_name", "schema":"default", "table": "test_flatten", "account_type": "ALIYUN", "access_id": "<yourAccessKeyId>", "access_key": "<yourAccessKeySecret>", "partition_window_type": "MINUTE", "mode":"VALUE", "format":"FLATTEN", "sink_pool_size":"150", "record_batch_size":"9000", "buffer_size_kb":"600000", "runtime.error.topic.name":"runtime_error", "runtime.error.topic.bootstrap.servers":"http://XXXX", "skip_error":"false" } } -
Execute o comando abaixo para iniciar a tarefa do conector Kafka.
curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" http://localhost:8083/connectors -d @$KAFKA_HOME/config/odps-sink-connector.json -
Verifique o resultado.
-
Consulte os dados da tabela do MaxCompute
Use um cliente local (odpscmd) ou outra ferramenta capaz de executar comandos SQL do MaxCompute para se conectar ao MaxCompute e, em seguida, execute o comando abaixo para consultar os dados e verificar o resultado.
set odps.sql.allow.fullscan=true; select * from table_flatten;A saída retornada é:
# As you can see from the results, the data with ID 101 was not written to MaxCompute because it did not match the table schema. # Because the runtime.error.topic.name parameter was configured, the process was not blocked, and subsequent data was written successfully. +-------+------------+------------+------------+------+------------+----+ | topic | partition | offset | id | name | extendinfo | pt | +-------+------------+------------+------------+------+------------+----+ | flatten_test | 0 | 0 | 123 | json-1 | {"gender":"M","region":"beijing"} | 07-14-2023 01:33 | | flatten_test | 0 | 1 | 456 | json-2 | {"gender":"W","region":"hangzhou"} | 07-14-2023 01:33 | | flatten_test | 0 | 0 | 123 | json-1 | {"gender":"M","region":"beijing"} | 07-14-2023 13:16 | | flatten_test | 0 | 1 | 456 | json-2 | {"gender":"W","region":"hangzhou"} | 07-14-2023 13:16 | | flatten_test | 0 | 2 | 100 | json-3 | {"gender":"M","region":"beijing"} | 07-14-2023 13:16 | | flatten_test | 0 | 4 | 102 | json-5 | {"gender":"M","region":"beijing"} | 07-14-2023 13:16 | +-------+------------+------------+------------+------+------------+----+ -
Consulte mensagens no tópico
runtime_errorNo diretório
$KAFKA_HOME/bin/, execute o comando abaixo para visualizar as mensagens.sh kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic runtime_error --from-beginningO resultado retornado é:
# The abnormal data is successfully written to the runtime_error message queue. {"id":101,"name":"json-4","extendinfos":"null"}
-