Este tópico descreve como criar um conector sink do MaxCompute para exportar dados de um tópico de origem de uma instância do ApsaraMQ for Kafka para uma tabela do MaxCompute.
Pré-requisitos
Atenda aos seguintes requisitos:
-
ApsaraMQ for Kafka
O recurso de conector está habilitado na instância do ApsaraMQ for Kafka. Para mais informações, consulte Enable Connector.
-
Crie um tópico na instância do ApsaraMQ for Kafka. Para mais informações, consulte Step 1: Create a topic.
Este exemplo utiliza um tópico chamado maxcompute-test-input.
-
MaxCompute
-
Crie uma tabela do MaxCompute no cliente do MaxCompute. Para mais informações, consulte Create a table.
Neste exemplo, uma tabela do MaxCompute chamada test_kafka é criada em um projeto chamado connector_test. Execute a seguinte instrução para criar essa tabela:
CREATE TABLE IF NOT EXISTS test_kafka(topic STRING,partition BIGINT,offset BIGINT,key STRING,value STRING) PARTITIONED by (pt STRING);
-
-
Opcional:EventBridge
O EventBridge está ativado. Para mais informações sobre como ativar o EventBridge, consulte Activate EventBridge and grant permissions.
NotaA ativação do EventBridge é necessária apenas quando a instância que contém o tópico de origem de dados está na região China (Hangzhou) ou China (Chengdu).
Precauções
Você só pode exportar dados de um tópico de origem de uma instância do ApsaraMQ for Kafka para uma tabela do MaxCompute dentro da mesma região. Para mais informações sobre os limites dos conectores, consulte Limits.
-
Se a instância com o tópico de origem estiver nas regiões China (Hangzhou) ou China (Chengdu), a tarefa do conector será publicada no EventBridge.
Atualmente, o EventBridge é gratuito. Para mais detalhes, consulte Billing.
-
Ao criar um conector, o EventBridge cria automaticamente a função vinculada ao serviço AliyunServiceRoleForEventBridgeSourceKafka.
Se a função vinculada ao serviço não existir, o EventBridge a criará automaticamente para permitir que o EventBridge acesse o ApsaraMQ for Kafka.
Se a função já estiver disponível, o EventBridge não criará uma nova.
Para mais informações sobre funções vinculadas ao serviço, consulte Service-linked roles for EventBridge.
Não é possível visualizar os logs operacionais das tarefas de conector publicadas no EventBridge. Após a conclusão da tarefa, verifique os detalhes de consumo dos grupos que assinam o tópico de origem para conferir o status. Para mais informações, consulte View consumer status.
Procedimento
Para exportar dados de um tópico de origem de uma instância do ApsaraMQ for Kafka para uma tabela do MaxCompute usando um conector sink do MaxCompute, siga as etapas abaixo:
-
Conceda permissões ao ApsaraMQ for Kafka para acessar o MaxCompute.
-
Opcional: Crie os tópicos e o grupo necessários para o conector sink do MaxCompute.
Se preferir não criar manualmente os tópicos e o grupo, pule esta etapa e defina o parâmetro Resource Creation Method como Auto na próxima etapa.
ImportanteTópicos específicos exigidos pelo conector sink do MaxCompute precisam usar um mecanismo de armazenamento local. Se a versão principal da sua instância do ApsaraMQ for Kafka for 0.10.2, esses tópicos não podem ser criados manualmente; a criação deve ser automática.
-
Verifique o resultado.
Criar uma função do RAM
Não é possível selecionar o ApsaraMQ for Kafka como serviço confiável durante a criação de uma função do RAM. Portanto, selecione inicialmente qualquer outro serviço confiável e depois modifique manualmente a política de confiança da função.
Faça login no RAM console.
No painel de navegação à esquerda, escolha .
Na página Roles, clique em Create Role.
-
No painel Create Role, execute as seguintes operações:
Selecione Alibaba Cloud Service como entidade confiável e clique em Next.
Defina o parâmetro Role Type como Normal Service Role. No campo RAM Role Name, insira AliyunKafkaMaxComputeUser1. Na lista suspensa Select Trusted Service, selecione MaxCompute. Em seguida, clique em OK.
Na página Roles, localize e clique em AliyunKafkaMaxComputeUser1.
Na página AliyunKafkaMaxComputeUser1, clique na aba Trust Policy Management e depois em Edit Trust Policy.
-
No painel Edit Trust Policy, substitua fc no script por alikafka e clique em OK.

Adicionar permissões
Para utilizar um conector sink do MaxCompute na exportação de mensagens para uma tabela do MaxCompute, conceda as seguintes permissões à função do RAM.
|
Objeto |
Operação |
Descrição |
|
Project |
CreateInstance |
Permissão para criar instâncias em projetos. |
|
Table |
Describe |
Permissão para ler metadados das tabelas. |
|
Table |
Alter |
Permissão para modificar metadados das tabelas, além de criar e excluir partições. |
|
Table |
Update |
Permissão para sobrescrever e inserir dados nas tabelas. |
Para mais detalhes sobre essas permissões e como concedê-las, consulte MaxCompute permissions.
Siga os passos abaixo para conceder as permissões necessárias à função AliyunKafkaMaxComputeUser1:
Faça login no cliente do MaxCompute.
-
Execute o comando a seguir para adicionar a função do RAM AliyunKafkaMaxComputeUser1 como usuário do RAM:
add user `RAM$<accountid>:role/aliyunkafkamaxcomputeuser1`;NotaSubstitua <accountid> pelo ID da sua conta Alibaba Cloud.
-
Conceda ao usuário do RAM as permissões mínimas necessárias para acessar o MaxCompute.
-
Execute o comando abaixo para conceder permissões ao usuário do RAM no projeto connector_test:
grant CreateInstance on project connector_test to user `RAM$<accountid>:role/aliyunkafkamaxcomputeuser1`;NotaSubstitua <accountid> pelo ID da sua conta Alibaba Cloud.
-
Execute o comando a seguir para conceder permissões ao usuário do RAM na tabela test_kafka:
grant Describe, Alter, Update on table test_kafka to user `RAM$<accountid>:role/aliyunkafkamaxcomputeuser1`;NotaSubstitua <accountid> pelo ID da sua conta Alibaba Cloud.
-
Criar os tópicos necessários para o conector sink do MaxCompute
No console do ApsaraMQ for Kafka, crie manualmente os cinco tópicos exigidos pelo conector sink do MaxCompute. São eles: tópico de offset da tarefa, tópico de configuração da tarefa, tópico de status da tarefa, tópico de fila de mensagens mortas e tópico de dados de erro. Esses tópicos diferem quanto ao número de partições e ao mecanismo de armazenamento. Para mais informações, consulte Parâmetros na etapa Configure Source Service.
Faça login no ApsaraMQ for Kafka console.
-
Na página Overview, selecione uma região na seção Resource Distribution.
ImportanteCrie os tópicos na mesma região da sua aplicação, ou seja, onde a instância ECS está implantada. Tópicos não funcionam entre regiões diferentes. Por exemplo, se criar um tópico na região China (Beijing), tanto o produtor quanto o consumidor de mensagens devem rodar em uma instância ECS também localizada em China (Beijing).
Na página Instances, clique no nome da instância desejada.
No painel de navegação à esquerda, clique em Topics.
Na página Topics, clique em Create Topic.
-
No painel Create Topic, configure as propriedades do tópico e clique em OK.
Parâmetro
Descrição
Exemplo
Name
Nome do tópico.
NotaNo Kafka, nomes de tópicos como
xxx_xxxexxx.xxxsão considerados idênticos. O sistema retornará um erro caso tente criar um tópico com nome duplicado.demo
Description
Breve descrição do tópico.
demo test
Partitions
Quantidade de partições no tópico.
12
Storage Engine
NotaAtualmente, a seleção do tipo de mecanismo de armazenamento está disponível apenas para instâncias non-Serverless Professional Edition. Para outras instâncias, este parâmetro não é suportado e assume o valor padrão Cloud Storage.
Mecanismo de armazenamento para as mensagens do tópico.
ApsaraMQ for Kafka suporta dois mecanismos de armazenamento.
Cloud Storage: A camada subjacente utiliza discos da Alibaba Cloud. Este mecanismo oferece baixa latência, alto desempenho, alta durabilidade e confiabilidade, empregando um mecanismo distribuído de três réplicas. Se a Instance Edition da instância for Standard (High Write), o único mecanismo disponível 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, a ordem das mensagens pode ser perdida. Ao definir Storage Engine como Cloud Storage, a opção Normal Message vem selecionada por padrão.
Partitionally Ordered Message: Por padrão, mensagens com a mesma chave vão para a mesma partição, mantendo a ordem de envio. Mesmo se houver falha em uma máquina do cluster, a ordem ainda é garantida dentro da partição. Contudo, o envio para algumas partições pode falhar temporariamente, normalizando após a recuperação. Quando Storage Engine for definido como Local Storage, a opção Partitionally Ordered Message será selecionada automaticamente.
Normal Message
Log Cleanup Policy
Política de limpeza de logs do tópico.
Ao configurar Storage Engine como Local Storage, é obrigatório definir a Log Cleanup Policy. Essa configuração só é permitida em instâncias Professional Edition; instâncias Standard Edition não a suportam.
ApsaraMQ for Kafka oferece duas políticas de limpeza.
Delete: Política padrão. Se houver espaço suficiente em disco, as mensagens permanecem pelo período máximo de retenção. Caso o disco fique insuficiente (geralmente acima de 85% de uso), mensagens antigas são excluídas antecipadamente para preservar a disponibilidade do serviço.
Compact: Aplica a política de Log Compaction do Kafka. Essa política garante que, para mensagens com a mesma chave, apenas o valor mais recente seja mantido. É ideal para cenários como recuperação de estado após falhas ou recarga de cache após reinicialização. Por exemplo, ao usar Kafka Connect ou Confluent Schema Registry, utilize um Kafka Compact Topic para armazenar estados do sistema ou configurações.
ImportanteTópicos compactados destinam-se geralmente a componentes específicos do ecossistema, como Kafka Connect ou Confluent Schema Registry. Evite aplicar essa propriedade em tópicos usados para outros cenários de envio e recebimento de mensagens. Para mais detalhes, consulte a Biblioteca de Demos do ApsaraMQ for Kafka.
Compact
Tag
Tags associadas ao tópico.
demo
Após a criação, o tópico aparecerá na lista da página Topics.
Criar o grupo necessário para o conector sink do MaxCompute
No console do ApsaraMQ for Kafka, crie manualmente o grupo exigido pelo conector sink do MaxCompute. O nome do grupo deve seguir o formato connect-nome da tarefa. Para mais detalhes, consulte Parâmetros na etapa Configure Source Service.
Faça login no ApsaraMQ for Kafka console.
Na página Overview, selecione uma região na seção Resource Distribution.
Na página Instances, clique no nome da instância desejada.
No painel de navegação à esquerda, clique em Groups.
Na página Groups, clique em Create Group.
-
No painel Create Group, insira o nome do grupo na caixa de texto Group ID, adicione uma breve descrição no campo Description, atribua tags ao grupo e clique em OK.
Após a criação, o grupo estará visível na lista da página Groups.
Criar e implantar um conector sink do MaxCompute
Para criar e implantar um conector sink do MaxCompute destinado à exportação de dados do ApsaraMQ for Kafka para o MaxCompute, execute os passos a seguir:
Faça login no ApsaraMQ for Kafka console.
Na página Overview, selecione uma região na seção Resource Distribution.
Na página Instances, clique no nome da instância desejada.
No painel de navegação à esquerda, clique em Connectors.
Na página Connectors, clique em Create Connector.
-
No assistente Create Connector, proceda da seguinte forma:
-
Na etapa Configure Basic Information, preencha os parâmetros descritos na tabela abaixo e clique em Next.
Parâmetro
Descrição
Exemplo
Name
Nome do conector. Observe as seguintes regras ao defini-lo:
Deve ter entre 1 e 48 caracteres, podendo conter dígitos, letras minúsculas e hífens (-), mas não pode iniciar com hífen.
Cada nome de conector deve ser exclusivo dentro de uma instância ApsaraMQ for Kafka.
O nome do grupo utilizado pela tarefa do conector deve seguir o formato connect-nome da tarefa. Caso ainda não exista tal grupo, o Message Queue for Apache Kafka o criará automaticamente.
kafka-maxcompute-sink
Instance
Informações sobre a instância do Message Queue for Apache Kafka. Por padrão, o nome e o ID da instância são exibidos.
demo alikafka_post-cn-st21p8vj****
-
Na etapa Configure Source Service, selecione Message Queue for Apache Kafka como serviço de origem, configure os parâmetros listados na tabela a seguir e clique em Next.
NotaSe você já criou um tópico e um consumer group, opte pela criação manual de recursos e informe os dados existentes. Caso contrário, escolha a criação automática.
Tabela 1. Parâmetros na etapa Configure Source Service
Parâmetro
Descrição
Exemplo
Data Source Topic
Nome do tópico de origem de dados de onde as informações serão exportadas.
maxcompute-test-input
Consumer Thread Concurrency
Número de threads de consumo concorrentes usadas para exportar dados do tópico de origem. Valor padrão: 6. Valores válidos:
1
2
3
6
12
6
Consumer Offset
Offset inicial para o consumo. Valores válidos:
Earliest Offset: O consumo inicia a partir do offset mais antigo.
Latest Offset: O consumo inicia a partir do offset mais recente.
Earliest Offset
VPC ID
ID da virtual private cloud (VPC) onde a tarefa de exportação será executada. Clique em Configure Runtime Environment para exibir o parâmetro. O valor padrão corresponde ao VPC ID especificado durante a implantação da instância ApsaraMQ for Kafka. Não é necessário alterá-lo.
vpc-bp1xpdnd3l***
vSwitch ID
ID do vSwitch onde a tarefa de exportação será executada. Clique em Configure Runtime Environment para visualizar o parâmetro. O vSwitch deve estar na mesma VPC da instância ApsaraMQ for Kafka. O valor padrão é o vSwitch ID informado na implantação da instância ApsaraMQ for Kafka.
vsw-bp1d2jgg81***
Failure Handling Policy
Define se a assinatura de uma partição deve ser mantida após falha no envio de mensagem. Clique em Configure Runtime Environment para ver o parâmetro. Valores válidos:
Continue Subscription: Mantém a assinatura na partição com erro e registra os logs.
Stop Subscription: Interrompe a assinatura na partição com erro e registra os logs.
NotaPara mais informações, consulte Manage connectors.
Para solucionar erros com base nos códigos de erro, consulte Error codes.
Continue Subscription
Resource Creation Method
Método para criar os tópicos e o grupo necessários ao conector sink do MaxCompute. Clique em Configure Runtime Environment para exibir o parâmetro.
Auto
Manual
Auto
Connector Consumer Group
Grupo utilizado pela tarefa de exportação de dados do conector. Clique em Configure Runtime Environment para visualizar o parâmetro. O nome do grupo deve seguir o padrão connect-nome da tarefa.
connect-kafka-maxcompute-sink
Task Offset Topic
Tópico usado para armazenar offsets de consumo. Clique em Configure Runtime Environment para exibir o parâmetro.
Topic: Recomenda-se iniciar o nome do tópico com connect-offset.
Partitions: A quantidade de partições deve ser maior que 1.
Storage Engine: Deve ser configurado como Local Storage.
cleanup.policy: A política de limpeza de logs deve ser definida como Compact.
connect-offset-kafka-maxcompute-sink
Task Configuration Topic
Tópico destinado a armazenar configurações da tarefa. Clique em Configure Runtime Environment para ver o parâmetro.
Topic: Recomenda-se iniciar o nome com connect-config.
Partitions: O tópico deve conter apenas uma partição.
Storage Engine: Configure como Local Storage.
cleanup.policy: Defina a política de limpeza como Compact.
connect-config-kafka-maxcompute-sink
Task Status Topic
Tópico utilizado para armazenar o status da tarefa. Clique em Configure Runtime Environment para mostrar o parâmetro.
Topic: Sugere-se começar o nome com connect-status.
Partitions: Recomenda-se definir 6 partições.
Storage Engine: Selecione Local Storage.
cleanup.policy: Escolha Compact como política de limpeza.
connect-status-kafka-maxcompute-sink
Dead-letter Queue Topic
Tópico para armazenar dados de erro do framework Kafka Connect. Clique em Configure Runtime Environment para acessar o parâmetro. Para economizar recursos, é possível usar o mesmo tópico como fila de mensagens mortas e error data topic.
Topic: Recomenda-se prefixar o nome com connect-error.
Partitions: Sugere-se configurar 6 partições.
Storage Engine: Pode ser Local Storage ou Cloud Storage.
connect-error-kafka-maxcompute-sink
Error Data Topic
Tópico para guardar dados de erro do conector. Clique em Configure Runtime Environment para visualizar o parâmetro. Visando economia de recursos, utilize o mesmo tópico para dead-letter queue topic e error data topic.
Topic: Recomenda-se iniciar o nome com connect-error.
Partitions: É aconselhável definir 6 partições.
Storage Engine: Aceita Local Storage ou Cloud Storage.
connect-error-kafka-maxcompute-sink
-
Na etapa Configure Destination Service, selecione MaxCompute como serviço de destino, ajuste os parâmetros conforme a tabela abaixo e clique em Create.
NotaSe a instância contendo o tópico de origem estiver nas regiões China (Hangzhou) ou China (Chengdu), a caixa de diálogo Service Authorization aparecerá ao selecionar MaxCompute como destino. Clique em OK na caixa Service Authorization, configure os parâmetros da tabela a seguir e clique em Create.
Parâmetro
Descrição
Exemplo
Endpoint
Endpoint do MaxCompute. Para mais detalhes, consulte Endpoint.
Endpoint de VPC: Recomendado por oferecer menor latência. Utilize-o quando a instância ApsaraMQ for Kafka e o projeto MaxCompute estiverem na mesma região.
Endpoint público: Não recomendado devido à maior latência. Use apenas se a instância ApsaraMQ for Kafka e o projeto MaxCompute estiverem em regiões distintas. Para utilizá-lo, habilite o acesso à Internet no conector. Consulte Enable Internet access for a connector.
http://service.cn-hangzhou.maxcompute.aliyun-inc.com/api
Workspace
Nome do projeto MaxCompute para onde os dados serão exportados.
connector_test
Table
Nome da tabela do MaxCompute que receberá os dados exportados.
test_kafka
Region for Table
Região onde a tabela do MaxCompute foi criada.
China (Hangzhou)
Alibaba Cloud Account ID
ID da conta Alibaba Cloud usada para acessar o MaxCompute.
188***
RAM Role
Nome da função do RAM assumida pelo ApsaraMQ for Kafka. Para mais informações, veja Criar uma função do RAM.
AliyunKafkaMaxComputeUser1
Mode
Modo de exportação das mensagens para o conector sink do MaxCompute. Valor padrão: DEFAULT. Valores válidos:
KEY: Apenas as chaves das mensagens são preservadas e gravadas na coluna Key da tabela do MaxCompute.
VALUE: Somente os valores das mensagens são mantidos e escritos na coluna Value da tabela do MaxCompute.
DEFAULT: Tanto chaves quanto valores são retidos. As chaves vão para a coluna Key e os valores para a coluna Value da tabela do MaxCompute.
ImportanteNo modo DEFAULT, o formato CSV não é suportado. Escolha apenas entre TEXT e BINARY.
DEFAULT
Format
Formato de exportação das mensagens para o conector sink do MaxCompute. Valor padrão: TEXT. Valores válidos:
TEXT: strings
BINARY: arrays de bytes
CSV: strings separadas por vírgulas (,)
ImportanteAo escolher CSV, o modo DEFAULT fica indisponível. Apenas KEY e VALUE são permitidos.
Modo KEY: Preserva apenas as chaves, separadas por vírgulas (,), gravando-as na tabela do MaxCompute seguindo a ordem dos índices.
Modo VALUE: Mantém somente os valores, separados por vírgulas (,), inserindo-os na tabela do MaxCompute conforme a ordem dos índices.
TEXT
Partition
Frequência de criação de partições. Valor padrão: HOUR. Valores válidos:
DAY: Grava dados em uma nova partição diariamente.
HOUR: Cria uma nova partição a cada hora.
MINUTE: Gera uma nova partição a cada minuto.
HOUR
Time Zone
Fuso horário do cliente produtor ApsaraMQ for Kafka que envia mensagens ao tópico de origem. Valor padrão: GMT 08:00.
GMT 08:00
Depois de criado, o conector ficará visível na página Connectors.
-
Acesse a página Connectors, localize o conector recém-criado e clique em Deploy na coluna Actions.
Enviar uma mensagem de teste
Após implantar o conector sink do MaxCompute, envie uma mensagem para o tópico de origem no ApsaraMQ for Kafka para validar se a exportação para o MaxCompute funciona corretamente.
Na página Connectors, encontre o conector desejado e clique em Test na coluna Actions.
-
No painel Send Message, envie uma mensagem de teste.
-
Defina Sending Method como Console.
Na caixa de texto Message Key, insira a chave da mensagem. Por exemplo, demo.
No campo Message Content, digite o conteúdo da mensagem de teste. Exemplo: {"key": "test"}.
-
Configure Send to Specified Partition para determinar se a mensagem irá para uma partição específica.
Clique em Yes e informe o ID da partição na caixa Partition ID. Por exemplo, 0. Para consultar IDs de partições, veja View partition status.
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.
Escolha Sending Method como SDK. Selecione um SDK compatível com a linguagem ou framework desejado e o tipo de conexão para enviar mensagens.
-
Visualizar dados na tabela do MaxCompute
Após enviar uma mensagem para o tópico de origem no ApsaraMQ for Kafka, faça login no cliente do MaxCompute para confirmar o recebimento.
Para visualizar a tabela test_kafka, siga estas etapas:
Faça login no cliente do MaxCompute.
-
Execute o comando abaixo para listar as partições da tabela:
show partitions test_kafka;Neste exemplo, o seguinte resultado é retornado:
pt=11-17-2020 15 OK -
Execute o comando a seguir para visualizar os dados armazenados nas partições:
select * from test_kafka where pt ="11-17-2020 14";Neste exemplo, o seguinte resultado é retornado:
+----------------------+------------+------------+-----+-------+---------------+ | topic | partition | offset | key | value | pt | +----------------------+------------+------------+-----+-------+---------------+ | maxcompute-test-input| 0 | 0 | 1 | 1 | 11-17-2020 14 | +----------------------+------------+------------+-----+-------+---------------+