Os triggers do ApsaraMQ for Kafka invocam funções do Function Compute quando mensagens são publicadas em um tópico do Kafka, o que permite o processamento orientado a eventos sem polling.
Como funciona
O ApsaraMQ for Kafka integra-se ao Function Compute por meio do EventBridge. Ao criar um trigger no console do Function Compute, o FC cria automaticamente fluxos de eventos no EventBridge com base na configuração do trigger.
Após a ativação, o trigger monitora o tópico especificado. Quando uma mensagem é publicada no ApsaraMQ for Kafka, o EventBridge entrega uma ou mais mensagens em lote à função associada, conforme as definições de lote.
Limites
A instância do ApsaraMQ for Kafka deve residir na mesma região da função.
Se o número de fluxos de eventos atingir o limite máximo, não será possível criar novos triggers do Kafka. Para consultar o limite, veja Limites.
Pré-requisitos
Antes de começar, verifique se você tem:
EventBridge: Ativou o EventBridge e concedeu permissões a um usuário RAM
Function Compute: Uma função de evento
ApsaraMQ for Kafka: Uma instância implantada e um tópico e grupo de consumidores
Etapa 1: Criar um trigger
Faça login no console do Function Compute e abra a função de destino.
Na aba Configurations, acesse a página Triggers e clique em Create Trigger.
Configure os parâmetros do trigger descritos na tabela a seguir e clique em OK.

|
Parâmetro |
Obrigatório |
Descrição |
Exemplo |
|
Consumer offset |
Sim |
Ponto a partir do qual o EventBridge começa a buscar mensagens. Earliest Offset: inicia a partir da mensagem disponível mais antiga. Latest Offset: inicia apenas a partir de novas mensagens. |
Earliest Offset |
|
Invocation method |
Sim |
Define como a função é invocada quando um evento ou lote de eventos é entregue. Sync Invocation: aguarda uma resposta antes de processar o próximo lote; limite de payload de 32 MB. Veja Invocação síncrona. Async Invocation: retorna imediatamente e continua para o próximo lote; limite de payload de 128 KB. Veja Invocação assíncrona. |
Sync Invocation |
|
Max. Delivery Concurrency |
Não |
Número máximo de mensagens do Kafka entregues simultaneamente ao Function Compute. Valores válidos: 1–300. Disponível apenas para Sync Invocation. Para aumentar o limite, acesse o Quota Center do EventBridge, localize EventStreaming FC Sink Maximum Concurrent Number of Synchronous Posting e clique em Apply. |
1 |
Para configurações avançadas, como definições de push, políticas de nova tentativa e filas de mensagens mortas, consulte Recursos avançados de triggers.
Etapa 2: (Opcional) Configurar parâmetros de teste
O trigger passa as mensagens do Kafka para a função como o parâmetro event. Para testar a função sem publicar uma mensagem real, configure manualmente um evento de teste que simule o payload do trigger.
Nota: O teste simulado é útil para validar sua lógica de parsing, mas não testa todo o pipeline do trigger. Para verificar o comportamento de ponta a ponta, use uma mensagem real do Kafka conforme descrito na Etapa 3.
Na aba Code, clique no ícone
ao lado de Test Function e selecione Configure Test Parameters.No painel Configure Test Parameters, clique em Create New Test Event ou Modify Existing Test Event, insira um nome e o conteúdo do evento e clique em OK.
O exemplo a seguir mostra a estrutura do event para duas mensagens em lote:
[
{
"specversion": "1.0",
"id": "8e215af8-ca18-4249-8645-f96c1026****",
"source": "acs:alikafka",
"type": "alikafka:Topic:Message",
"subject": "acs:alikafka_pre-cn-i7m2t7t1****:topic:mytopic",
"datacontenttype": "application/json; charset=utf-8",
"time": "2022-06-23T02:49:51.589Z",
"aliyunaccountid": "164901546557****",
"data": {
"topic": "****",
"partition": 7,
"offset": 25,
"timestamp": 1655952591589,
"headers": {
"headers": [],
"isReadOnly": false
},
"key": "keytest",
"value": "hello kafka msg"
}
},
{
"specversion": "1.0",
"id": "8e215af8-ca18-4249-8645-f96c1026****",
"source": "acs:alikafka",
"type": "alikafka:Topic:Message",
"subject": "acs:alikafka_pre-cn-i7m2t7t1****:topic:mytopic",
"datacontenttype": "application/json; charset=utf-8",
"time": "2022-06-23T02:49:51.589Z",
"aliyunaccountid": "164901546557****",
"data": {
"topic": "****",
"partition": 7,
"offset": 25,
"timestamp": 1655952591589,
"headers": {
"headers": [],
"isReadOnly": false
},
"key": "keytest",
"value": "hello kafka msg"
}
}
]
O objeto data contém os campos específicos do Kafka:
|
Campo |
Tipo |
Exemplo |
Descrição |
|
|
String |
TopicName |
Nome do tópico. |
|
|
Int |
1 |
Número da partição de onde a mensagem foi consumida. |
|
|
Int |
0 |
Offset da mensagem dentro da partição. |
|
|
String |
1655952591589 |
Timestamp Unix (ms) que indica quando o consumo da mensagem começou. |
Para os campos de envelope CloudEvents (specversion, id, source, type, etc.), consulte Visão geral de eventos.
Etapa 3: Escrever e testar o código da função
-
Na aba Code, escreva o código da função no editor e clique em Deploy. O exemplo Node.js a seguir mostra o handler da função:
'use strict'; /* To enable the initializer feature please implement the initializer function as below: exports.initializer = (context, callback) => { console.log('initializing'); callback(null, ''); }; */ exports.handler = (event, context, callback) => { console.log("event: %s", event); // Parse the event parameters and process the event. callback(null, 'return result'); } -
Teste a função usando um dos métodos a seguir:
Teste simulado: Clique em Test Function para usar o evento de teste configurado na Etapa 2. Esse método é útil para validar sua lógica de parsing sem uma conexão ativa com o Kafka.
Teste de ponta a ponta: Faça login no console do ApsaraMQ for Kafka, selecione o tópico e clique em Send Message para publicar uma mensagem real. O trigger detecta a mensagem e invoca a função automaticamente.

-
Após a execução, visualize a saída em Real-time Logs.

Próximos passos
Para modificar ou excluir um trigger existente, consulte Gerenciar triggers.