Todos os produtos
Search
Central de documentação

MaxCompute:,

Última atualização: Jun 27, 2026

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

  1. Este exemplo utiliza um ambiente Linux. Em uma janela de comando, baixe o pacote kafka-connector-2.0.jar executando o comando abaixo ou usando o link de download.

    wget http://maxcompute-repo.oss-cn-hangzhou.aliyuncs.com/kafka/kafka-connector-2.0.jar

    Para evitar conflitos de dependência, crie uma subpasta, como connector, no diretório $KAFKA_HOME/libs e coloque o pacote kafka-connector-2.0.jar dentro dela.

    Nota

    Caso o pacote kafka-connector-2.0.jar nã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ço Kafka-connector.

  2. No diretório $KAFKA_HOME/config, configure o arquivo connect-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  
  3. No diretório $KAFKA_HOME/, execute o comando abaixo para iniciar o serviço Kafka-connector.

    ## Start command
    bin/connect-distributed.sh config/connect-distributed.properties &

Configure e inicie a tarefa do conector Kafka

  1. Crie e configure o arquivo de configuração odps-sink-connector.json. Em seguida, carregue o arquivo odps-sink-connector.json em 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.

        Nota
        • Se 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.

  2. 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

  1. 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 usa topic_text como nome do tópico.

      sh kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic topic_text

      Execute 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
  2. (Opcional) Inicie o serviço Kafka-connector. Para mais informações, consulte Configure e inicie o serviço do conector Kafka.

    Nota

    Se o serviço Kafka-connector já estiver em execução, pule esta etapa.

  3. Crie e configure o arquivo odps-sink-connector.json. Em seguida, carregue o arquivo odps-sink-connector.json em 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 arquivo odps-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"
        }
      }
  4. 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
  5. 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

  1. 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 chamado topic_csv.

      sh kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic topic_csv

      Execute 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
  2. (Opcional) Inicie o serviço Kafka-connector. Para mais informações, consulte Configure e inicie o serviço do conector Kafka.

    Nota

    Se o serviço Kafka-connector já estiver em execução, pule esta etapa.

  3. Crie e configure o arquivo odps-sink-connector.json e, em seguida, carregue o arquivo odps-sink-connector.json em qualquer local. Este tópico usa o caminho $KAFKA_HOME/config como exemplo.

    O código abaixo fornece um exemplo do arquivo odps-sink-connector.json. Para mais informações sobre o arquivo odps-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"
        }
      }
    
  4. 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
  5. 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

  1. 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 usa topic_json como nome do tópico.

      sh kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic topic_json

      Execute 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"}
  2. (Opcional) Inicie o serviço Kafka-connector. Para mais informações, consulte Configure e inicie o serviço do conector Kafka.

    Nota

    Se o serviço Kafka-connector já estiver em execução, pule esta etapa.

  3. Crie e configure o arquivo odps-sink-connector.json. Em seguida, carregue o arquivo odps-sink-connector.json em 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 arquivo odps-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"
        }
      }
    
  4. 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
  5. 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

  1. 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 usa topic_flatten como nome do tópico.

      ./kafka/bin/kafka-topics.sh --create --bootstrap-server localhost:9092 --replication-factor 1 --partitions 1 --topic topic_flatten

      Execute 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"}}
  2. (Opcional) Inicie o serviço Kafka-connector. Para mais informações, consulte Configure e inicie o serviço do conector Kafka.

    Nota

    Se o serviço Kafka-connector já estiver em execução, pule esta etapa.

  3. Crie e configure o arquivo odps-sink-connector.json e, em seguida, carregue o arquivo odps-sink-connector.json em qualquer local. Este tópico usa o caminho $KAFKA_HOME/config como exemplo.

    O código abaixo fornece um exemplo do arquivo odps-sink-connector.json. Para mais informações sobre o arquivo odps-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"
        }
      }
    
  4. 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
  5. 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

  1. 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_error
        Nota

        Se 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"}} 
  2. (Opcional) Inicie o serviço Kafka-connector. Para mais informações, consulte Configure e inicie o serviço do conector Kafka.

    Nota

    Se o serviço Kafka-connector já estiver em execução, pule esta etapa.

  3. Crie e configure o arquivo odps-sink-connector.json e, em seguida, carregue o arquivo odps-sink-connector.json em qualquer local. Este tópico usa o caminho $KAFKA_HOME/config como exemplo.

    O código abaixo fornece um exemplo do arquivo odps-sink-connector.json. Para mais informações sobre o arquivo odps-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"
      }
    }
    
  4. 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
  5. 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_error

      No 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-beginning

      O resultado retornado é:

      # The abnormal data is successfully written to the runtime_error message queue.
      {"id":101,"name":"json-4","extendinfos":"null"}