Todos os produtos
Search
Central de documentação

ApsaraMQ for Kafka:Criar um conector sink do Tablestore

Última atualização: Jun 27, 2026

Este tópico descreve como criar um conector sink do Tablestore para exportar dados de um tópico em uma instância do ApsaraMQ for Kafka para o Tablestore.

Pré-requisitos

Observações

  • É possível exportar dados de um tópico de origem em uma instância do ApsaraMQ for Kafka para o Tablestore apenas dentro da mesma região. Para obter detalhes sobre os limites do conector, consulte Limites de uso.

  • Ao criar um conector, o ApsaraMQ for Kafka cria automaticamente uma função vinculada ao serviço.

    • Se a função vinculada ao serviço não existir, o ApsaraMQ for Kafka a criará automaticamente, permitindo a exportação de dados para o Tablestore.

    • Caso a função vinculada ao serviço já exista, o ApsaraMQ for Kafka não criará uma nova.

    Para mais informações sobre funções vinculadas ao serviço, consulte Funções vinculadas ao serviço.

Procedimento

Este tópico descreve como usar um conector sink do Tablestore para exportar dados de um tópico de origem em uma instância do ApsaraMQ for Kafka para o Tablestore.

  1. Opcional: Crie os tópicos e o grupo necessários para o conector sink do Tablestore.

    Se não for necessário personalizar os tópicos e o grupo, pule esta etapa e selecione Auto na próxima.

    Importante

    Alguns tópicos exigidos pelo conector sink do Tablestore devem usar o mecanismo de armazenamento Local. Em instâncias do ApsaraMQ for Kafka com versão principal 0.10.2, não é possível criar manualmente tópicos com o mecanismo de armazenamento Local. A criação desses tópicos ocorre apenas de forma automática.

    1. Criar os tópicos para o conector sink do Tablestore

    2. Criar o grupo para o conector sink do Tablestore

  2. Criar e implantar o conector sink do Tablestore

  3. Verifique o resultado

    1. Enviar uma mensagem de teste

    2. Visualizar dados na tabela

Criar tópicos para o conector sink do Tablestore

No console do ApsaraMQ for Kafka, crie manualmente os cinco tópicos necessários para o conector sink do Tablestore: tópico de offset de tarefa, tópico de configuração de tarefa, tópico de status de tarefa, tópico de fila de mensagens mortas e tópico de dados de erro. O número necessário de partições e o mecanismo de armazenamento variam conforme o tópico. Para mais informações, consulte Lista de parâmetros do serviço de origem.

  1. Faça login no console do ApsaraMQ for Kafka.

  2. Na página Overview, selecione uma região na seção Resource Distribution.

    Importante

    Crie os tópicos na mesma região da sua aplicação, ou seja, onde a instância ECS está implantada. Não é possível usar tópicos entre regiões diferentes. Por exemplo, se você criar um tópico na região China (Beijing), o produtor e o consumidor de mensagens também deverão ser executados em uma instância ECS na região China (Beijing).

  3. Na página Instances, clique em nome da instância desejada.

  4. No painel de navegação à esquerda, clique em Topics.

  5. Na página Topics, clique em Create Topic.

  6. No painel Create Topic, defina as propriedades do tópico e clique em OK.

    Parâmetro

    Descrição

    Exemplo

    Name

    Nome do tópico.

    Nota

    No Kafka, nomes de tópicos como xxx_xxx e xxx.xxx são considerados iguais. O sistema reportará um erro se você tentar criar um tópico com nome duplicado.

    demo

    Description

    Breve descrição do tópico.

    demo test

    Partitions

    Número de partições no tópico.

    12

    Storage Engine

    Nota

    Atualmente, a seleção do tipo de mecanismo de armazenamento está disponível apenas para instâncias da Professional Edition não Serverless. Para outras instâncias, este parâmetro não é suportado e assume o valor padrão Cloud Storage.

    Mecanismo de armazenamento das mensagens do tópico.

    ApsaraMQ for Kafka suporta os dois mecanismos de armazenamento a seguir.

    • Cloud Storage: A camada subjacente utiliza discos do Alibaba Cloud. Este mecanismo oferece baixa latência, alto desempenho, alta durabilidade e alta confiabilidade por meio de um mecanismo distribuído de três réplicas. Se a Instance Edition da instância for Standard (High Write), o único mecanismo de armazenamento permitido será Cloud Storage.

    • Local Storage: Utiliza o algoritmo nativo de replicação ISR (in-sync replica) do Kafka e um mecanismo distribuído de três réplicas.

    Cloud Storage

    Message Type

    Tipo das mensagens do tópico.

    • Normal Message: Por padrão, mensagens com a mesma chave são distribuídas para a mesma partição, e as mensagens dentro de uma partição são armazenadas na ordem de envio. Se uma máquina no cluster falhar, as mensagens poderão ficar desordenadas. Ao definir o Storage Engine como Cloud Storage, a opção Normal Message é selecionada por padrão.

    • Partitionally Ordered Message: Por padrão, mensagens com a mesma chave são distribuídas para a mesma partição, e as mensagens dentro de uma partição são armazenadas na ordem de envio. Se uma máquina no cluster falhar, a ordem das mensagens ainda será garantida dentro da partição. No entanto, o envio de mensagens para algumas partições pode falhar. As operações normais serão retomadas após a recuperação das partições. Ao definir o Storage Engine como Local Storage, a opção Partitionally Ordered Message é selecionada por padrão.

    Normal Message

    Log Cleanup Policy

    Política de limpeza dos logs do tópico.

    Ao definir o Storage Engine como Local Storage, configure obrigatoriamente a Log Cleanup Policy. A definição do Storage Engine como Local Storage é suportada apenas em instâncias da Professional Edition. Instâncias da Standard Edition não são suportadas.

    ApsaraMQ for Kafka suporta as duas políticas de limpeza a seguir.

    • Delete: Política padrão de limpeza de mensagens. Se houver capacidade de disco suficiente, as mensagens serão retidas pelo período máximo de retenção. Caso contrário (geralmente quando o uso excede 85%), mensagens antigas serão excluídas antecipadamente para garantir a disponibilidade do serviço.

    • Compact: Utiliza a política de limpeza Log Compaction do Kafka. Essa política garante que, para mensagens com a mesma chave, o valor mais recente seja sempre retido. Aplica-se principalmente a cenários como recuperação de estado após falha do sistema ou recarregamento de cache após reinicialização. Por exemplo, ao usar Kafka Connect ou Confluent Schema Registry, use um Kafka Compact Topic para armazenar o estado do sistema ou informações de configuração.

      Importante

      Tópicos compactados geralmente são usados apenas em componentes específicos do ecossistema, como Kafka Connect ou Confluent Schema Registry. Não defina esta propriedade para tópicos em outros cenários de envio e recebimento de mensagens. Para mais informações, consulte Biblioteca de Demos do ApsaraMQ for Kafka.

    Compact

    Tag

    Tags do tópico.

    demo

    Após a criação, o tópico aparecerá na lista da página Topics.

Criar um grupo para o conector sink do Tablestore

Crie manualmente um grupo para a tarefa de sincronização de dados do conector sink do Tablestore no console do ApsaraMQ for Kafka. O nome deste grupo deve ser connect-Nome da Tarefa. Para mais informações, consulte Parâmetros do serviço de origem.

  1. Faça login no console do ApsaraMQ for Kafka.

  2. Na página Overview, selecione uma região na seção Resource Distribution.

  3. Na página Instances, clique em nome da instância desejada.

  4. No painel de navegação à esquerda, clique em Groups.

  5. Na página Groups, clique em Create Group.

  6. No painel Create Group, insira um nome para o grupo na caixa de texto Group ID, insira uma breve descrição na caixa de texto Description, adicione tags ao grupo e clique em OK.

    Após a criação, o grupo aparecerá na lista da página Groups.

Criar e implantar um conector sink do Tablestore

Crie e implante um conector sink do Tablestore para sincronizar dados de uma instância do ApsaraMQ for Kafka para uma tabela do Tablestore.

  1. Faça login no console do ApsaraMQ for Kafka.

  2. Na página Overview, selecione uma região na seção Resource Distribution.

  3. No painel de navegação à esquerda, clique em Connectors.

  4. Na página Connectors, selecione a instância à qual o conector pertence na lista suspensa Select Instance e clique em Create Connector.

  5. No assistente Create Connector, conclua as etapas a seguir.

    1. Na aba Configure Basic Information, configure os parâmetros a seguir e clique em Next.

      Parâmetro

      Descrição

      Exemplo

      Name

      Nome do conector. O nome deve atender aos seguintes requisitos:

      • Pode ter até 48 caracteres e conter apenas dígitos, letras minúsculas e hifens (-). Não pode começar com hífen (-).

      • Deve ser único dentro de uma instância do ApsaraMQ for Kafka.

      A tarefa de sincronização de dados do conector usa um grupo de consumidores chamado connect-task-name. Se você não criar este grupo de consumidores manualmente, o sistema o criará automaticamente.

      kafka-ts-sink

      Instance

      Por padrão, o nome e o ID da instância são exibidos.

      demo alikafka_post-cn-st21p8vj****

    2. Na aba Configure Source Service, defina Data Source como ApsaraMQ for Kafka, configure os parâmetros a seguir e clique em Next.

      Nota

      Se você já criou o tópico e o grupo necessários, selecione Manual para a criação de recursos e insira as informações correspondentes. Caso contrário, selecione Auto.

      Tabela 1. Parâmetros do serviço de origem

      Parâmetro

      Descrição

      Exemplo

      Data Source Topic

      Tópico de origem para a sincronização de dados.

      ts-test-input

      Consumer Thread Concurrency

      Número de threads de consumo concorrentes para o tópico de origem de dados. O valor padrão é 6. Valores válidos:

      • 1

      • 2

      • 3

      • 6

      • 12

      6

      Consumer Offset

      Posição inicial para o consumo de mensagens. Valores válidos:

      • Earliest Offset: Consome mensagens a partir do offset mais antigo.

      • Latest Offset: Consome mensagens a partir do offset mais recente.

      Earliest Offset

      VPC ID

      ID da Virtual Private Cloud (VPC) onde a tarefa de sincronização de dados é executada. Este parâmetro aparece após clicar em Configure Runtime Environment. O valor padrão é a VPC da instância do ApsaraMQ for Kafka e não requer configuração.

      vpc-bp1xpdnd3l***

      vSwitch ID

      ID do vSwitch onde a tarefa de sincronização de dados é executada. Este parâmetro aparece após clicar em Configure Runtime Environment. O vSwitch deve estar na mesma VPC da instância do ApsaraMQ for Kafka. O valor padrão é o vSwitch especificado durante a implantação da instância do ApsaraMQ for Kafka.

      vsw-bp1d2jgg81***

      Failure Handling Policy

      Política para lidar com falhas no envio de mensagens em uma partição de tópico. Este parâmetro aparece após clicar em Configure Runtime Environment. Valores válidos:

      • Continue Subscription: Continua a assinatura na partição do tópico onde ocorreu o erro e imprime logs de erro.

      • Stop Subscription: Interrompe a assinatura na partição do tópico onde ocorreu o erro e imprime logs de erro.

      Nota

      Continue Subscription

      Resource Creation Method

      Método para criar os tópicos e o grupo necessários. Este parâmetro aparece após clicar em Configure Runtime Environment.

      • Auto

      • Manual

      Auto

      Connector Consumer Group

      Grupo de consumidores usado pela tarefa de sincronização de dados. Este parâmetro aparece após clicar em Configure Runtime Environment. O nome do grupo deve seguir o formato connect-task-name.

      connect-cluster-kafka-ots-sink

      Task Offset Topic

      Tópico que armazena os offsets dos consumidores. Este parâmetro aparece após clicar em Configure Runtime Environment.

      • Tópico: O nome do tópico deve começar com connect-offset.

      • Partições: O número de partições do tópico deve ser maior que 1.

      • Mecanismo de armazenamento: O mecanismo de armazenamento do tópico deve ser Local Storage.

      • cleanup.policy: A política de limpeza de log do tópico deve ser compact.

      connect-offset-kafka-ots-sink

      Task Configuration Topic

      Tópico que armazena as configurações da tarefa. Este parâmetro aparece após clicar em Configure Runtime Environment.

      • Tópico: O nome do tópico deve começar com connect-config.

      • Partições: O número de partições do tópico deve ser 1.

      • Mecanismo de armazenamento: O mecanismo de armazenamento do tópico deve ser Local Storage.

      • cleanup.policy: A política de limpeza de log do tópico deve ser compact.

      connect-config-kafka-ots-sink

      Task Status Topic

      Tópico que armazena o status da tarefa. Este parâmetro aparece após clicar em Configure Runtime Environment.

      • Tópico: O nome do tópico deve começar com connect-status.

      • Partições: O número recomendado de partições é 6.

      • Mecanismo de armazenamento: O mecanismo de armazenamento do tópico deve ser Local Storage.

      • cleanup.policy: A política de limpeza de log do tópico deve ser compact.

      connect-status-kafka-ots-sink

      Dead-letter Queue Topic

      Tópico usado para armazenar dados de erro do framework Kafka Connect. Este parâmetro aparece após clicar em Configure Runtime Environment. Para economizar recursos, use o mesmo tópico para a fila de mensagens mortas e o Error Data Topic.

      • Tópico: O nome do tópico deve começar com connect-error.

      • Partições: O número recomendado de partições é 6.

      • Mecanismo de armazenamento: O mecanismo de armazenamento do tópico pode ser Local Storage ou Cloud Storage.

      connect-error-kafka-ots-sink

      Error Data Topic

      Tópico usado para armazenar dados de erro do sink. Este parâmetro aparece após clicar em Configure Runtime Environment. Para economizar recursos, use o mesmo tópico para o Dead-letter Queue Topic e o tópico de dados de erro.

      • Tópico: O nome do tópico deve começar com connect-error.

      • Partições: O número recomendado de partições é 6.

      • Mecanismo de armazenamento: O mecanismo de armazenamento do tópico pode ser Local Storage ou Cloud Storage.

      connect-error-kafka-ots-sink

    3. Na aba Configure Destination Service, defina Destination Service como Tablestore, configure os parâmetros a seguir e clique em Create.

      Parâmetro

      Descrição

      Exemplo

      Instance Name

      Nome da instância do Tablestore.

      k00eny67****

      Automatically Create Destination Table

      Define se uma tabela deve ser criada automaticamente no Tablestore.

      • Yes: Uma tabela é criada automaticamente no Tablestore para armazenar os dados sincronizados, com base no nome da tabela configurado.

      • No: Uma tabela existente é usada para armazenar os dados sincronizados.

      Yes

      Destination Table Name

      Nome da tabela que armazena os dados sincronizados. Se você definir Automatically Create Destination Table como No, o nome da tabela deve ser igual ao de uma tabela existente na instância do Tablestore.

      kafka_table

      Tablestore

      Tipo da tabela que armazena os dados sincronizados.

      • Wide Column Model

      • TimeSeries Model

      Wide Column Model

      Message Key Format

      Formato de entrada da chave da mensagem. Os valores válidos são String e JSON. O valor padrão é JSON. Este parâmetro aparece apenas quando Tablestore está definido como Wide Column Model.

      • String: A chave da mensagem é analisada diretamente como string.

      • JSON: A chave da mensagem deve estar no formato JSON.

      String

      Message Value Format

      Formato de entrada do valor da mensagem. Os valores válidos são String e JSON. O valor padrão é JSON. Este parâmetro aparece apenas quando Tablestore está definido como Wide Column Model.

      • String: O valor da mensagem é analisado diretamente como string.

      • JSON: O valor da mensagem deve estar no formato JSON.

      String

      JSON Message Field Conversion

      Método para processar campos em uma mensagem JSON. Este parâmetro aparece se você definir Message Key Format ou Message Value Format como JSON. Valores válidos:

      • Write All as String: Converte todos os campos para o tipo String no Tablestore.

      • Automatically Identify Field Types: Converte campos String e Boolean no corpo da mensagem JSON para os tipos String e Boolean correspondentes no Tablestore. Os tipos Integer e Float no corpo da mensagem JSON são convertidos para o tipo Double no Tablestore.

      Write All as String

      Primary Key Mode

      Especifica o modo de chave primária. É possível extrair chaves primárias da tabela de diferentes partes dos registros de mensagem do ApsaraMQ for Kafka, incluindo coordenadas (tópico, partição e offset), chave e valor. Este parâmetro aparece apenas quando Tablestore está definido como Wide Column Model. O valor padrão é kafka.

      • kafka: Usa <connect_topic>_<connect_partition> e <connect_offset> como chaves primárias da tabela de dados.

      • record_key: Usa campos na chave do registro como chaves primárias da tabela de dados.

      • record_value: Usa campos no valor do registro como chaves primárias da tabela de dados.

      kafka

      Primary Key Column Names

      Nomes das colunas de chave primária da tabela de dados e seus tipos de dados correspondentes. Os tipos String e Integer são suportados. Isso significa que os campos na chave ou valor do registro correspondentes aos nomes de coluna configurados serão usados como chaves primárias da tabela de dados.

      Este parâmetro aparece se você definir Message Key Format como JSON e Primary Key Mode como record_key, ou se definir Message Value Format como JSON e Primary Key Mode como record_value.

      Clique em Create para adicionar um nome de coluna. É possível configurar no máximo quatro nomes de coluna.

      None

      Write Mode

      Especifica o modo de escrita. Os valores válidos são put e update. O valor padrão é put. Este parâmetro aparece apenas quando Tablestore está definido como Wide Column Model.

      • put: Sobrescreve os dados existentes.

      • update: Atualiza os dados existentes.

      put

      Delete Mode

      Se um registro de mensagem do ApsaraMQ for Kafka contiver valores nulos, escolha entre excluir linhas ou colunas de atributo. Este parâmetro aparece quando Primary Key Mode está definido como record_key. Valores válidos:

      • none: Valor padrão. Nenhuma exclusão é permitida.

      • row: Permite a exclusão de linhas.

      • column: Permite a exclusão de colunas de atributo.

      • row_and_column: Permite a exclusão de linhas e colunas de atributo.

      O comportamento de exclusão depende do modo de escrita:

      • Se o Write Mode for put, qualquer modo de exclusão resultará em sobrescrita na tabela de dados do Tablestore, mesmo quando o valor contiver campos nulos.

      • Se o Write Mode for update e o modo de exclusão for none ou row, registros com todos os campos de valor nulos serão tratados como dados incorretos. Se apenas alguns campos de valor forem nulos, o conector ignorará os campos nulos e gravará os campos não nulos na tabela de dados do Tablestore. Se o modo de exclusão for column ou row_and_column, o conector excluirá linhas e colunas de atributo para campos nulos e, em seguida, gravará os dados na tabela de dados do Tablestore.

      None

      Metric Name Field

      Campo mapeado para o campo de nome da métrica (_m_name) no TimeSeries Model do Tablestore. O nome da métrica especifica a grandeza física ou a métrica de monitoramento dos dados de série temporal, como temperatura ou velocidade. Este campo não pode estar vazio. O parâmetro aparece quando você seleciona o TimeSeries Model para Tablestore.

      measurement

      Data Source Field

      Mapeie este campo para o campo de origem de dados (_data_source) no TimeSeries Model do Tablestore. Este campo atua como identificador da origem dos dados de série temporal, como nome de máquina ou ID de dispositivo, e pode estar vazio. O parâmetro aparece quando o modelo TimeSeries é selecionado para Tablestore.

      source

      Tag Field

      Use um ou mais campos como campos de tag (_tags) para um TimeSeries Model do Tablestore. Cada tag é um par chave-valor de strings. A chave é o nome do campo configurado e o valor é o conteúdo do campo. As tags fazem parte dos metadados da série temporal. Uma série temporal é identificada exclusivamente pela combinação de seu nome de métrica, origem de dados e tags. As tags podem estar vazias. O parâmetro aparece quando Tablestore é selecionado como modelo de série temporal.

      tag1, tag2

      Timestamp Field

      Mapeia este campo para o campo de carimbo de data/hora (_time) no TimeSeries Model do Tablestore. Representa o ponto no tempo desta linha de dados de série temporal, como o momento em que uma grandeza física é gerada. Quando os dados são gravados no Tablestore, o campo de carimbo de data/hora é convertido em microssegundos para gravação e armazenamento. O parâmetro aparece quando você seleciona o TimeSeries Model para Tablestore.

      time

      Timestamp Unit

      Configure este parâmetro com base no campo de carimbo de data/hora real. O parâmetro aparece apenas quando Tablestore está definido como TimeSeries Model. Valores válidos:

      • SECONDS

      • MILLISECONDS

      • MICROSECONDS

      • NANOSECONDS

      MILLISECONDS

      Whether to Map All Non-primary Key Fields

      Define se todos os campos que não são de chave primária devem ser mapeados como campos de dados. Campos que não são de chave primária são aqueles ainda não mapeados como nome de métrica, origem de dados, tag ou carimbo de data/hora. O parâmetro aparece apenas quando Tablestore está definido como TimeSeries Model. Valores válidos:

      • Yes: Os campos são mapeados automaticamente e seus tipos de dados são determinados. Tipos numéricos são convertidos para o tipo Double.

      • No: Especifique os campos e tipos a serem mapeados.

      Yes

      Configure Mapping for All Non-primary Key Fields

      Tipos de campo correspondentes aos nomes de campos que não são de chave primária da tabela de séries temporais. Os tipos Double, Integer, String, Binary e Boolean são suportados. O parâmetro aparece se você definir Whether to Map All Non-primary Key Fields como No.

      String

      Após a criação, visualize o conector na página Connectors.

  6. Após criar o conector, localize-o na página Connectors e clique em Deploy na coluna Actions.

  7. Clique em OK.

Enviar uma mensagem de teste

Após implantar o conector sink do Tablestore, envie uma mensagem para o tópico de origem no ApsaraMQ for Kafka para verificar se os dados estão sendo sincronizados com o Tablestore.

  1. Na página Connectors, localize o conector desejado e clique em Test na coluna Actions.

  2. No painel Send Message, envie uma mensagem de teste.

    • Defina Sending Method como Console.

      1. Na caixa de texto Message Key, insira a chave da mensagem. Por exemplo, demo.

      2. Na caixa de texto Message Content, insira o conteúdo da mensagem de teste. Por exemplo, {"key": "test"}.

      3. Defina Send to Specified Partition para especificar se a mensagem deve ser enviada para uma partição específica.

        • Clique em Yes e insira o ID da partição na caixa de texto Partition ID. Por exemplo, 0. Para consultar o ID da partição, consulte Visualizar status da partição.

        • Clique em No para não especificar uma partição.

    • Defina Sending Method como Docker. Execute o comando Docker na seção Run the Docker container to produce a sample message para enviar a mensagem.

    • Defina Sending Method como SDK. Selecione um SDK para a linguagem ou framework necessário e um tipo de conexão para enviar mensagens.

Visualizar dados da tabela

Após enviar uma mensagem para o tópico de origem de dados no ApsaraMQ for Kafka, visualize os dados na tabela do Tablestore para confirmar o recebimento da mensagem.

  1. Faça login no console do Tablestore.

  2. Na página Overview, clique em nome da instância ou em Instances na coluna Actions.

  3. Na aba Instance Details, localize a tabela desejada na seção Tables.

  4. Clique em nome da tabela. Na página Table Manage, clique na aba Data Management para visualizar os dados.