Todos os produtos
Search
Central de documentação

DataWorks:Fonte de dados Kafka

Última atualização: Jun 29, 2026

A fonte de dados Kafka oferece um canal bidirecional para leitura e gravação de dados no Kafka. Este tópico descreve as capacidades de sincronização de dados que o DataWorks fornece para o Kafka.

Versões suportadas

O DataWorks suporta o Alibaba Cloud Kafka e versões auto-gerenciadas do Kafka, da 0.10.2 até a 3.6.x.

Nota

A sincronização de dados não é suportada para versões do Kafka anteriores à 0.10.2. Isso ocorre porque essas versões não permitem recuperar offsets de partição e suas estruturas de dados podem não suportar timestamps.

Leitura em tempo real

  • Ao utilizar um grupo de recursos Serverless por assinatura, estime as especificações necessárias antecipadamente para evitar falhas nas tarefas devido à insuficiência de recursos.

    Estime 1 CU por tópico. Também é necessário estimar os recursos com base no tráfego:

    • Para dados do Kafka sem compressão, estime 1 CU para cada 10 MB/s de tráfego.

    • Para dados do Kafka compactados, estime 2 CU para cada 10 MB/s de tráfego.

    • Caso os dados compactados exijam análise JSON, estime 3 CU para cada 10 MB/s de tráfego.

  • Ao usar um grupo de recursos Serverless por assinatura ou uma versão antiga de um grupo de recursos exclusivo para Integração de Dados:

    • Se sua carga de trabalho tiver alta tolerância a failover, o uso de slots do cluster não deve exceder 80%.

    • Para cargas de trabalho com baixa tolerância a failover, mantenha o uso de slots do cluster abaixo de 70%.

Nota

O consumo real de recursos depende de fatores como conteúdo e formato dos dados. Após a avaliação inicial, ajuste os recursos conforme o uso observado durante a execução.

Limitações

A fonte de dados Kafka suporta grupos de recursos Serverless (recomendado) e versões antigas de grupos de recursos exclusivos para Integração de Dados.

Leitura offline de tabela única

Se tanto parameter.groupId quanto parameter.kafkaConfig.group.id estiverem configurados, parameter.groupId terá precedência sobre group.id no parâmetro kafkaConfig.

Gravação em tempo real em tabela única

Operações de gravação não suportam deduplicação de dados. Se uma tarefa for reiniciada após uma redefinição de offset ou um failover, dados duplicados poderão ser gravados.

Gravação em tempo real para banco de dados inteiro

  • Tarefas de sincronização de dados em tempo real são compatíveis com grupos de recursos Serverless (recomendado) e versões antigas de grupos de recursos exclusivos para Integração de Dados.

  • Quando a tabela de origem possui chave primária, o valor dessa chave é usado como chave do registro Kafka. Isso garante que alterações na mesma chave primária sejam gravadas na mesma partição do Kafka em ordem.

  • Se a tabela de origem não tiver chave primária, há duas opções. Ao escolher sincronizar tabelas sem chave primária, a chave do registro Kafka ficará vazia. Para garantir que as alterações da tabela sejam gravadas no Kafka em ordem, o tópico Kafka de destino deve ter apenas uma única partição. Caso opte por uma chave primária personalizada, uma combinação de um ou mais campos que não sejam chave primária será usada como chave do registro Kafka.

  • Para assegurar que alterações da mesma chave primária sejam gravadas ordenadamente na mesma partição do Kafka, mesmo se o cluster Kafka retornar uma exceção, adicione a seguinte configuração no formulário de parâmetros estendidos.

    {"max.in.flight.requests.per.connection":1,"buffer.memory": 100554432}

    Importante

    Essa configuração reduz significativamente o desempenho de replicação. É preciso equilibrar o desempenho com a necessidade de ordenação estrita e confiabilidade.

  • Para mais detalhes sobre o formato geral das mensagens gravadas no Kafka em sincronizações em tempo real, o formato das mensagens de heartbeat e o formato das mensagens correspondentes a alterações nos dados de origem, consulte Apêndice: Formato de mensagem.

Tipos de campo suportados

O Kafka oferece armazenamento de dados não estruturados. Um registro Kafka geralmente inclui os seguintes campos: key, value, offset, timestamp, headers e partition. Ao ler ou gravar dados no Kafka, o DataWorks processa as informações da seguinte forma.

Leitura de dados

Durante a leitura de dados do Kafka, o DataWorks pode analisar as informações no formato JSON. A tabela abaixo descreve como cada módulo de dados é processado.

Módulo de dados do registro Kafka

Tipo de dados processado

key

Depende do item de configuração keyType na tarefa de sincronização de dados. Para mais informações sobre o parâmetro keyType, consulte a descrição completa dos parâmetros no apêndice.

value

Determinado pelo item de configuração valueType na tarefa de sincronização de dados. Consulte a descrição completa dos parâmetros no apêndice para detalhes sobre valueType.

offset

Long

timestamp

Long

headers

String

partition

Long

Gravação de dados

O DataWorks grava dados no Kafka em formato JSON ou texto. A política de processamento varia conforme o tipo de tarefa de sincronização, conforme descrito na tabela abaixo.

Importante
  • Na gravação em formato de texto, os nomes dos campos não são incluídos. Os valores dos campos são separados por um delimitador.

  • Em tarefas de sincronização em tempo real que gravam no Kafka, utiliza-se o formato JSON integrado. Os dados incluem informações como mensagens de alteração de banco de dados, horário comercial e informações de DDL (Data Definition Language). Para detalhes sobre o formato dos dados, consulte Apêndice: Formato de mensagem.

Tipo de tarefa de sincronização

Formato do value gravado no Kafka

Tipo de campo de origem

Método de processamento para operações de gravação

Sincronização offline

Nó de sincronização offline no DataStudio

JSON

String

String codificada em UTF-8

Boolean

Convertido para a string UTF-8 "true" ou "false"

Time/Date

String codificada em UTF-8 no formato yyyy-MM-dd HH:mm:ss

Numérico

String numérica codificada em UTF-8

Fluxo de bytes

O fluxo de bytes é tratado como uma string codificada em UTF-8 e convertido para string.

text

String

String codificada em UTF-8

Boolean

Convertido para a string UTF-8 "true" ou "false"

Time/Date

String codificada em UTF-8 no formato yyyy-MM-dd HH:mm:ss

Numérico

String numérica codificada em UTF-8

Fluxo de bytes

O fluxo de bytes é tratado como uma string codificada em UTF-8 e convertido para string.

Sincronização em tempo real: ETL em tempo real para Kafka

Nó de sincronização em tempo real no DataStudio

JSON

String

String codificada em UTF-8

Boolean

Tipo Boolean JSON

Time/Date

  • Para valores de tempo com precisão inferior a milissegundos: Convertido para um inteiro JSON de 13 dígitos representando o timestamp em milissegundos.

  • Para valores de tempo com precisão de microssegundos ou nanossegundos: Convertido para um número de ponto flutuante JSON contendo um inteiro de 13 dígitos para o timestamp em milissegundos e uma parte decimal de 6 dígitos para o timestamp em nanossegundos.

Numérico

Tipo numérico JSON

Fluxo de bytes

O fluxo de bytes é codificado em Base64 e depois convertido para uma string codificada em UTF-8.

text

String

String codificada em UTF-8

Boolean

Convertido para a string UTF-8 "true" ou "false"

Time/Date

String codificada em UTF-8 no formato yyyy-MM-dd HH:mm:ss

Numérico

String numérica codificada em UTF-8

Fluxo de bytes

O fluxo de bytes é codificado em Base64 e depois convertido para uma string codificada em UTF-8.

Sincronização em tempo real: Sincronização em tempo real de banco de dados inteiro para Kafka

Sincronização em tempo real apenas de dados incrementais

Formato JSON integrado

String

String codificada em UTF-8

Boolean

Tipo Boolean JSON

Time/Date

Timestamp de milissegundos com 13 dígitos

Numérico

Valor numérico JSON

Fluxo de bytes

O fluxo de bytes é codificado em Base64 e depois convertido para uma string codificada em UTF-8.

Solução de sincronização: Sincronização em tempo real com um clique para Kafka

Sincronização offline completa + sincronização incremental em tempo real

Formato JSON integrado

String

String codificada em UTF-8

Boolean

Tipo Boolean JSON

Time/Date

Timestamp de milissegundos com 13 dígitos

Numérico

Valor numérico JSON

Fluxo de bytes

O fluxo de bytes é codificado em Base64 e depois convertido para uma string codificada em UTF-8.

Adicionar uma fonte de dados

Antes de desenvolver uma tarefa de sincronização no DataWorks, adicione a fonte de dados necessária seguindo as instruções em Gerenciamento de fontes de dados. Consulte as descrições de parâmetros no console do DataWorks para entender o significado de cada parâmetro ao adicionar uma fonte de dados.

Desenvolver uma tarefa de sincronização de dados

Para informações sobre o ponto de entrada e o procedimento de configuração de uma tarefa de sincronização, consulte os guias de configuração a seguir.

Configure uma tarefa de sincronização offline para tabela única

Configure uma tarefa de sincronização em tempo real para tabela única ou banco de dados inteiro

Para instruções, consulte Configure uma tarefa de sincronização em tempo real para tabela única e Configure uma tarefa de sincronização em tempo real para banco de dados inteiro.

Configuração de autenticação

SSL

Ao configurar uma fonte de dados Kafka, se você definir o Special Authentication Method como SSL ou SASL_SSL, a autenticação SSL será ativada para o cluster Kafka. É obrigatório enviar o arquivo de certificado truststore do cliente e inserir a senha do truststore.

O arquivo de certificado keystore, a senha do keystore e a senha SSL são necessários apenas quando a autenticação SSL bidirecional está habilitada no cluster Kafka. O servidor do cluster Kafka utiliza esses dados para autenticar a identidade do cliente. A autenticação SSL bidirecional é ativada quando ssl.client.auth=required está definido no arquivo server.properties do cluster Kafka. Para mais informações, consulte Usar criptografia SSL para conexões Kafka.

GSSAPI

Se você definir o Sasl Mechanism como GSSAPI ao configurar uma fonte de dados Kafka, será necessário enviar três arquivos de autenticação: um arquivo de configuração JAAS, um arquivo de configuração Kerberos e um arquivo Keytab. Também é preciso configurar as definições de DNS/HOST para o grupo de recursos exclusivo. As seções a seguir detalham esses arquivos e as configurações necessárias de DNS e HOST.

Nota

Para um grupo de recursos Serverless, configure as informações de endereço do host usando resolução DNS interna. Para mais detalhes, consulte Resolução DNS interna (PrivateZone).

  • Arquivo de configuração JAAS

    O arquivo JAAS deve começar com KafkaClient, seguido por todos os itens de configuração entre chaves {}:

    • A primeira linha dentro das chaves define a classe do componente de logon a ser usada. Para diferentes mecanismos de autenticação SASL, essa classe é fixa. Cada item subsequente segue o formato key=value.

    • Todos os itens de configuração, exceto o último, não devem terminar com ponto e vírgula.

    • O último item de configuração deve terminar com ponto e vírgula, e outro ponto e vírgula deve seguir a chave de fechamento }.

    Se os requisitos de formato não forem atendidos, o arquivo de configuração JAAS não poderá ser analisado. O código abaixo mostra um formato típico de arquivo JAAS. Substitua os marcadores xxx pelas suas informações reais.

    KafkaClient {
       com.sun.security.auth.module.Krb5LoginModule required
       useKeyTab=true
       keyTab="xxx"
       storeKey=true
       serviceName="kafka-server"
       principal="kafka-client@EXAMPLE.COM";
    };

    Item de configuração

    Descrição

    Módulo de logon

    Deve ser definido como com.sun.security.auth.module.Krb5LoginModule.

    useKeyTab

    Defina como true.

    keyTab

    Aceita qualquer caminho. Durante a execução da tarefa de sincronização, o sistema baixa automaticamente o arquivo keytab enviado na configuração da fonte de dados para um caminho local, utilizando-o neste item de configuração.

    storeKey

    Indica se o cliente salva a chave. Pode ser true ou false, sem impacto na sincronização de dados.

    serviceName

    Corresponde ao item sasl.kerberos.service.name no arquivo server.properties do servidor Kafka. Configure conforme necessário.

    principal

    Principal Kerberos usado pelo cliente Kafka. Configure conforme necessário e certifique-se de que o arquivo keytab enviado contenha a chave para este principal.

  • Arquivo de configuração Kerberos

    O arquivo de configuração Kerberos deve conter dois módulos: [libdefaults] e [realms].

    • O módulo [libdefaults] especifica os parâmetros de autenticação Kerberos. Cada item segue o formato key=value.

    • O módulo [realms] define o endereço do KDC (Key Distribution Center). Pode conter vários submódulos realm, cada um iniciando com o nome do realm seguido de um sinal de igual (=).

    Em seguida, vem um conjunto de itens de configuração entre chaves, também no formato key=value. O código abaixo ilustra um formato típico de arquivo Kerberos. Substitua os marcadores xxx pelas suas informações reais.

    [libdefaults]
      default_realm = xxx
    
    [realms]
      xxx = {
        kdc = xxx
      }

    Item de configuração

    Descrição

    [libdefaults].default_realm

    Realm padrão usado ao acessar nós do cluster Kafka. Geralmente coincide com o realm do principal do cliente especificado no arquivo JAAS.

    Outros parâmetros de [libdefaults]

    O módulo [libdefaults] pode definir outros parâmetros de autenticação Kerberos, como ticket_lifetime. Configure conforme necessário.

    Nome do realm em [realms]

    Deve corresponder ao realm do principal do cliente no arquivo JAAS e ao [libdefaults].default_realm. Se houver diferença entre eles, inclua dois submódulos realms, correspondentes respectivamente ao realm do principal no JAAS e ao [libdefaults].default_realm.

    [realms].nome_do_realm.kdc

    Especifica o endereço e porta do KDC no formato IP:port, por exemplo, kdc=10.0.0.1:88. Se a porta for omitida, o sistema usa a porta padrão 88, como em kdc=10.0.0.1.

  • Arquivo Keytab

    O arquivo keytab deve conter a chave do principal especificado no arquivo JAAS e ser verificável pelo KDC. Por exemplo, se houver um arquivo chamado client.keytab no diretório de trabalho atual, execute o comando abaixo para verificar se ele contém a chave do principal desejado.

    klist -ket ./client.keytab
    
    Keytab name: FILE:client.keytab
    KVNO Timestamp           Principal
    ---- ------------------- ------------------------------------------------------
       7 2018-07-30T10:19:16 te**@**.com (des-cbc-md5)
  • Configuração de DNS e HOST para grupo de recursos exclusivo

    Quando um cluster Kafka usa autenticação Kerberos, o KDC registra o principal de cada nó usando seu hostname. Ao se conectar a um nó do cluster, o cliente usa as configurações locais de DNS e HOST para derivar o principal do nó e solicitar uma credencial de acesso ao KDC. Para acessar um cluster Kafka com Kerberos habilitado através de um grupo de recursos exclusivo, configure corretamente o DNS e o HOST para garantir a obtenção das credenciais de acesso aos nós do cluster via KDC:

    • Configurações de DNS

      Se você utilizar uma instância PrivateZone para resolução de nomes dos nós do cluster Kafka na VPC à qual o grupo de recursos exclusivo está conectado, adicione uma rota personalizada para os IPs 100.100.2.136 e 100.100.2.138 na conexão VPC. Isso garante que as configurações de resolução de nomes do PrivateZone para os nós do Kafka se apliquem ao grupo de recursos exclusivo. No painel de navegação à esquerda do console DataWorks, clique em resource group list. Na coluna Actions do seu grupo de recursos exclusivo, clique em network settings. Na aba VPC attachment, clique em custom route na coluna Actions. Na caixa de diálogo exibida, clique em Add Route, defina Destination Type como IDC e Connection Method como Direct IP, insira o endereço IP direto e clique em Generate Route.

    • Configurações de HOST

      Caso não utilize uma instância PrivateZone para resolver os nomes dos nós do cluster Kafka na VPC conectada ao grupo de recursos exclusivo, adicione os mapeamentos de IP para domínio de cada nó na configuração de host. No painel de navegação à esquerda do console DataWorks, clique em resource group list. Localize seu grupo de recursos, clique em network settings na coluna Actions e selecione a aba host configuration. Clique em Add para incluir um domínio de host. A configuração de host tem precedência sobre as configurações de DNS.

PLAIN

Ao configurar uma fonte de dados Kafka, se o Sasl Mechanism for definido como PLAIN, o arquivo JAAS deve iniciar com KafkaClient, seguido pelos itens de configuração entre chaves {}.

  • A primeira linha dentro das chaves define a classe do componente de logon. Para diferentes mecanismos SASL, essa classe é fixa. Os itens subsequentes seguem o formato key=value.

  • Exceto pelo último item, nenhum item de configuração deve terminar com ponto e vírgula.

  • O último item deve terminar com ponto e vírgula. Adicione também um ponto e vírgula após a chave de fechamento "}".

Se o formato não estiver correto, o arquivo JAAS não será interpretado. O código abaixo apresenta um modelo típico. Substitua os placeholders xxx pelos seus dados reais.

KafkaClient {
  org.apache.kafka.common.security.plain.PlainLoginModule required
  username="xxx"
  password="xxx";
};

Item de configuração

Descrição

Módulo de logon

Deve ser definido como org.apache.kafka.common.security.plain.PlainLoginModul

username

Nome de usuário. Configure conforme necessário.

password

Senha. Configure conforme necessário.

Perguntas frequentes

Apêndice: Exemplos de scripts e descrições de parâmetros

Configure uma tarefa de sincronização em lote usando o editor de código

Para configurar uma tarefa de sincronização em lote via editor de código, defina os parâmetros relevantes no script seguindo os requisitos unificados de formato. Para mais detalhes, consulte Configuração em modo script. As informações a seguir descrevem os parâmetros obrigatórios para fontes de dados nesse tipo de configuração.

Exemplo de script do Reader

A configuração JSON abaixo lê dados do Kafka.

{
    "type": "job",
    "steps": [
        {
            "stepType": "kafka",
            "parameter": {
                "server": "host:9093",
                "column": [
                    "__key__",
                    "__value__",
                    "__partition__",
                    "__offset__",
                    "__timestamp__",
                    "'123'",
                    "event_id",
                    "tag.desc"
                ],
                "kafkaConfig": {
                    "group.id": "demo_test"
                },
                "topic": "topicName",
                "keyType": "ByteArray",
                "valueType": "ByteArray",
                "beginDateTime": "20190416000000",
                "endDateTime": "20190416000006",
                "skipExceedRecord": "true"
            },
            "name": "Reader",
            "category": "reader"
        },
        {
            "stepType": "stream",
            "parameter": {
                "print": false,
                "fieldDelimiter": ","
            },
            "name": "Writer",
            "category": "writer"
        }
    ],
    "version": "2.0",
    "order": {
        "hops": [
            {
                "from": "Reader",
                "to": "Writer"
            }
        ]
    },
    "setting": {
        "errorLimit": {
            "record": "0"
        },
        "speed": {
            "throttle": true,//When throttle is false, the mbps parameter does not take effect, which means that the rate is not limited. When throttle is true, the rate is limited.
            "concurrent": 1,//The number of concurrent threads.
            "mbps":"12"//The maximum transmission rate. 1 mbps is equal to 1 MB/s.
        }
    }
}

Parâmetros do script do Reader

Parâmetro

Descrição

Obrigatório

datasource

Nome da fonte de dados. O editor de código permite adicionar fontes de dados. O valor deste parâmetro deve corresponder exatamente ao nome da fonte adicionada.

Sim

server

Endereço do servidor broker Kafka no formato IP:port.

É possível configurar apenas um server, mas garanta que o DataWorks consiga se conectar aos IPs de todos os brokers no cluster Kafka.

Sim

topic

Tópico Kafka. Representa uma agregação de fluxos de mensagens processados pelo Kafka.

Sim

column

Dados do Kafka a serem lidos. Suporta colunas constantes, colunas de dados e colunas de atributos.

  • Coluna constante: Coluna entre aspas simples, como ["'abc'", "'123'"].

  • Coluna de dados

    • Para dados em formato JSON, é possível extrair propriedades do objeto, como ["event_id"].

    • Para dados JSON, também é possível acessar subpropriedades aninhadas, como ["tag.desc"].

  • Coluna de atributo

    • __key__: Chave da mensagem.

    • __value__: Conteúdo completo da mensagem.

    • __partition__: Partição onde a mensagem atual reside.

    • __headers__: Cabeçalhos da mensagem atual.

    • __offset__: Offset da mensagem atual.

    • __timestamp__: Timestamp da mensagem atual.

    O código abaixo fornece um exemplo completo.

    "column": [
        "__key__",
        "__value__",
        "__partition__",
        "__offset__",
        "__timestamp__",
        "'123'",
        "event_id",
        "tag.desc"
        ]

Sim

keyType

Tipo da chave Kafka. Valores válidos: BYTEARRAY, DOUBLE, FLOAT, INTEGER, LONG e SHORT.

Não

valueType

Tipo do valor Kafka. Valores válidos: BYTEARRAY, DOUBLE, FLOAT, INTEGER, LONG e SHORT.

Não

beginDateTime

Horário inicial para consumo de dados. Define o limite esquerdo (inclusivo) do intervalo de tempo. É uma string no formato yyyymmddhhmmss. Compatível com scheduling parameters. Para mais informações, consulte Formatos suportados de parâmetros de agendamento.

Nota

Recurso disponível a partir do Kafka 0.10.2.

Especifique este parâmetro ou beginOffset.

Nota

beginDateTime e endDateTime são usados em conjunto.

endDateTime

Horário final para consumo de dados. Define o limite direito (exclusivo) do intervalo de tempo. É uma string no formato yyyymmddhhmmss. Compatível com scheduling parameters. Para mais informações, consulte Formatos suportados de parâmetros de agendamento.

Nota

Recurso disponível a partir do Kafka 0.10.2.

Especifique este parâmetro ou endOffset.

Nota

endDateTime e beginDateTime são usados em conjunto.

beginOffset

Offset inicial para consumo de dados. Pode ser configurado das seguintes formas:

  • Um número, como 15553274, indicando o offset inicial de consumo.

  • seekToBeginning: Consome dados a partir do offset mais antigo.

  • seekToLast: Lê dados a partir do offset salvo para o group ID definido em group.id no parâmetro kafkaConfig. Note que o offset do grupo é enviado automaticamente ao servidor Kafka pelo cliente em intervalos regulares. Portanto, se uma tarefa falhar e for reexecutada, pode ocorrer duplicação ou perda de dados. Se o parâmetro skipExceedRecord estiver definido como true, a tarefa poderá descartar os últimos registros lidos. Como o offset desses dados já foi confirmado no servidor, eles não serão lidos na próxima execução.

  • seekToEnd: Consome dados a partir do offset mais recente. Resultará em leitura de dados vazios.

Especifique este parâmetro ou beginDateTime.

endOffset

Offset final para consumo de dados. Controla quando a tarefa de consumo deve encerrar.

Especifique este parâmetro ou endDateTime.

skipExceedRecord

O Kafka usa public ConsumerRecords<K, V> poll(final Duration timeout) para consumir dados. Uma única chamada poll pode buscar dados além do endOffset ou endDateTime especificado. Este parâmetro determina se esses dados excedentes serão gravados no destino. Como a tarefa usa confirmação automática de offset, recomenda-se:

  • Para versões do Kafka anteriores à 0.10.2: Defina skipExceedRecord como false.

  • Para Kafka 0.10.2 e posteriores: Defina skipExceedRecord como true.

Não. O valor padrão é false.

partition

Um tópico Kafka possui múltiplas partições (partition). Por padrão, a tarefa lê dados de um intervalo de offsets cobrindo todas as partições do tópico. Também é possível especificar uma partition para ler dados apenas do intervalo de offsets de uma única partição.

Não. Sem valor padrão.

kafkaConfig

Ao criar um cliente KafkaConsumer para consumo de dados, é possível definir parâmetros estendidos como bootstrap.servers, auto.commit.interval.ms e session.timeout.ms. Use kafkaConfig para controlar o comportamento de consumo do KafkaConsumer.

Não

encoding

Quando keyType ou valueType é STRING, a codificação definida aqui é usada para interpretar a string.

Não. O valor padrão é UTF-8.

waitTIme

Tempo máximo, em segundos, que o consumidor aguarda para puxar dados do Kafka em uma única tentativa.

Não. O valor padrão é 60.

stopWhenPollEmpty

Valores válidos: true e false. Se definido como true e o consumidor receber dados vazios do Kafka (geralmente porque todos os dados do tópico foram lidos ou devido a problemas de rede/disponibilidade), a tarefa para imediatamente. Caso contrário, tenta novamente até obter dados.

Não. O valor padrão é true.

stopWhenReachEndOffset

Este parâmetro só tem efeito quando stopWhenPollEmpty é true. Valores válidos: true e false.

  • Se definido como true e o consumidor não receber dados numa requisição poll, verifica se o offset mais recente das partições foi atingido. Se todas as partições estiverem no offset mais recente, a tarefa para. Caso contrário, continua tentando puxar dados.

  • Se definido como false e não houver dados no poll, a verificação não ocorre e a tarefa para imediatamente.

Não. O valor padrão é false.

Nota

Configuração para compatibilidade retroativa. Versões do Kafka anteriores à 0.10.2 não suportam verificação do offset mais recente de todas as partições.

A tabela a seguir descreve os parâmetros kafkaConfig.

Parâmetro

Descrição

fetch.min.bytes

Quantidade mínima de dados, em bytes, que o consumidor busca do broker em uma única requisição. O broker aguarda até ter essa quantidade disponível antes de responder.

fetch.max.wait.ms

Tempo máximo, em milissegundos, que o broker aguarda dados ficarem disponíveis antes de responder a uma requisição fetch. Padrão: 500. O broker responde quando fetch.min.bytes ou fetch.max.wait.ms for satisfeito.

max.partition.fetch.bytes

Define o número máximo de bytes que o broker pode retornar ao consumidor por partition. Padrão: 1 MB.

session.timeout.ms

Tempo que o consumidor pode ficar desconectado do servidor antes de parar de receber serviços. Padrão: 30 segundos.

auto.offset.reset

Ação tomada pelo consumidor ao ler sem offset ou com offset inválido (devido a inatividade prolongada e expiração do registro). Padrão: none (sem redefinição automática). Altere para earliest para ler registros da partition desde o offset mais antigo.

max.poll.records

Número de mensagens retornadas por chamada ao método poll.

key.deserializer

Método de desserialização da chave da mensagem, como org.apache.kafka.common.serialization.StringDeserializer.

value.deserializer

Método de desserialização do valor dos dados, como org.apache.kafka.common.serialization.StringDeserializer.

ssl.truststore.location

Caminho do certificado raiz SSL.

ssl.truststore.password

Senha do armazenamento de certificados raiz. Para Alibaba Cloud Kafka, use KafkaOnsClient.

security.protocol

Protocolo de acesso. Atualmente, apenas SASL_SSL é suportado.

sasl.mechanism

Método de autenticação SASL. Para Alibaba Cloud Kafka, utilize PLAIN.

java.security.auth.login.config

Caminho do arquivo de autenticação SASL.

Exemplo de script do Writer

O código abaixo mostra a configuração JSON para gravar dados no Kafka.

{
  "type":"job",
  "version":"2.0",//The version number.
  "steps":[
    {
      "stepType":"stream",
      "parameter":{},
      "name":"Reader",
      "category":"reader"
    },
    {
      "stepType":"Kafka",//The plug-in name.
      "parameter":{
          "server": "ip:9092", //The server address of Kafka.
          "keyIndex": 0, //The column to be used as the key. Must follow camel case naming conventions, with k in lowercase.
          "valueIndex": 1, //The column to be used as the value. Currently, you can only select one column from the source data or leave this parameter empty. If empty, all source data is used.
          //For example, to use the 2nd, 3rd, and 4th columns of an ODPS table as the kafkaValue, create a new ODPS table, clean and integrate the data from the original ODPS table into the new table, and then use the new table for synchronization.
          "keyType": "Integer", //The type of the Kafka key.
          "valueType": "Short", //The type of the Kafka value.
          "topic": "t08", //The Kafka topic.
          "batchSize": 1024 //The amount of data written to Kafka at a time, in bytes.
        },
      "name":"Writer",
      "category":"writer"
    }
  ],
  "setting":{
      "errorLimit":{
      "record":"0"//The number of error records.
    },
    "speed":{
        "throttle":true,//When throttle is false, the mbps parameter does not take effect, which means that the rate is not limited. When throttle is true, the rate is limited.
        "concurrent":1, //The number of concurrent jobs.
        "mbps":"12"//The maximum transmission rate. 1 mbps is equal to 1 MB/s.
    }
   },
    "order":{
    "hops":[
        {
            "from":"Reader",
            "to":"Writer"
        }
      ]
    }
}

Parâmetros do script do Writer

Parâmetro

Descrição

Obrigatório

datasource

Nome da fonte de dados. O editor de código permite adicionar fontes de dados. O valor deve corresponder exatamente ao nome da fonte adicionada.

Sim

server

Endereço do servidor Kafka no formato IP:port.

Sim

topic

Tópico Kafka. Categoria para diferentes fluxos de mensagens processados pelo Kafka.

Cada mensagem publicada em um cluster Kafka pertence a uma categoria chamada tópico, que agrupa um conjunto de mensagens.

Sim

valueIndex

Coluna no writer Kafka usada como valor. Se não especificado, todas as colunas são concatenadas para formar o valor, usando o separador definido em fieldDelimiter.

Não

writeMode

Quando valueIndex não está configurado, este parâmetro define o formato de concatenação das colunas do registro de origem para formar o valor do registro Kafka. Valores válidos: text e JSON. Padrão: text.

  • Em modo text, as colunas são concatenadas usando o separador de fieldDelimiter.

  • Em modo JSON, as colunas formam uma string JSON baseada nos nomes definidos no parâmetro column.

Por exemplo, se um registro de origem tem três colunas com valores a, b e c, writeMode é text e fieldDelimiter é #, o valor gravado será a string a#b#c. Se writeMode for JSON e column for [{"name":"col1"},{"name":"col2"},{"name":"col3"}], o valor será {"col1":"a","col2":"b","col3":"c"}.

Se valueIndex estiver configurado, este parâmetro é ignorado.

Não

column

Campos da tabela de destino onde os dados serão gravados, separados por vírgulas. Exemplo: "column": ["id", "name", "age"].

Quando valueIndex não está configurado e writeMode é JSON, este parâmetro define os nomes dos campos na estrutura JSON para os valores das colunas de origem. Exemplo: "column": [{"name":id","type":"JSON_NUMBER"}, {"name":"name","type":"JSON_STRING"}, {"name":"age","type":"JSON_NUMBER"}].

  • Se o número de colunas na origem for maior que o configurado em column, os dados são truncados. Exemplo:

    Se a origem tem três colunas (a, b, c) e column é [{"name":"col1","type":"JSON_STRING"},{"name":"col2","type":"JSON_STRING"}], o valor gravado será {"col1":"a","col2":"b"}.

  • Se a origem tiver menos colunas que o especificado em column, os campos extras recebem null ou a string definida em nullValueFormat. Exemplo:

    Se a origem tem duas colunas (a, b) e column é [{"name":"col1","type":"JSON_STRING"},{"name":"col2","type":"JSON_STRING"},{"name":"col3","type":"JSON_STRING"}], o resultado será {"col1":"a","col2":"b","col3":null}. Se valueIndex estiver configurado ou writeMode for text, este parâmetro é ignorado.

  • Se o tipo de campo JSON não for definido, assume-se JSON_STRING.

  • Para valores válidos de tipos de campo JSON, consulte Apêndice: Tipos de campo JSON.

Se valueIndex estiver configurado ou writeMode for text, este parâmetro é ignorado.

Obrigatório quando valueIndex não está configurado e writeMode é JSON.

partition

Define o número da partição no tópico Kafka onde os dados serão gravados. Deve ser um inteiro maior ou igual a 0.

Não

keyIndex

Coluna no writer Kafka usada como chave.

O valor de keyIndex deve ser um inteiro maior ou igual a 0. Caso contrário, a tarefa falhará.

Não

keyIndexes

Array com os números ordinais das colunas de origem usadas como chave do registro Kafka.

A numeração começa em 0. Por exemplo, [0,1,2] concatena os valores dessas colunas com vírgulas para formar a chave. Se não especificado, a chave será null e os dados serão distribuídos em round-robin pelas partições. Use apenas este parâmetro ou keyIndex.

Não

fieldDelimiter

Quando writeMode é text e valueIndex não está configurado, as colunas de origem são concatenadas usando este separador para formar o valor do registro Kafka. Aceita caracteres únicos ou múltiplos, incluindo Unicode no formato \u0001. Caracteres de escape como \t e \n são suportados. Padrão: \t.

Ignorado se writeMode não for text ou se valueIndex estiver configurado.

Não

keyType

Tipo da chave Kafka. Valores válidos: BYTEARRAY, DOUBLE, FLOAT, INTEGER, LONG e SHORT.

Sim

valueType

Tipo do valor Kafka. Valores válidos: BYTEARRAY, DOUBLE, FLOAT, INTEGER, LONG e SHORT.

Sim

nullKeyFormat

Se o valor da coluna de origem definida por keyIndex ou keyIndexes for null, ele será substituído pela string especificada aqui. Se não configurado, nenhuma substituição ocorre.

Não

nullValueFormat

Se o valor de uma coluna de origem for null, ele será substituído pela string definida aqui ao montar o valor do registro Kafka. Sem substituição se não especificado.

Não

acks

Configuração acks ao inicializar o producer Kafka. Determina o método de confirmação de gravações bem-sucedidas. Por padrão, acks é all. Valores válidos para acks:

  • 0: Sem confirmação de gravação.

  • 1: Confirmação de gravação na réplica primária.

  • all: Confirmação de gravação em todas as réplicas.

Não

Apêndice: Definição de formato de mensagem para gravação no Kafka

Após configurar e executar uma tarefa de sincronização em tempo real, os dados lidos do banco de origem são gravados em um tópico Kafka no formato JSON. Inicialmente, todos os dados existentes na tabela de origem são transferidos para o tópico correspondente. Em seguida, a sincronização em tempo real começa a escrever continuamente dados incrementais. Informações de alterações DDL incrementais da tabela de origem também são gravadas no formato JSON. É possível obter status e informações de alterações das mensagens no Kafka. Para mais detalhes, consulte Apêndice: Formato de mensagem.

Nota

Nos dados JSON de uma tarefa de sincronização offline, os campos payload.sequenceId, payload.timestamp.eventTime e payload.timestamp.checkpointTime são definidos como -1.

Apêndice: Tipos de campo JSON

Quando writeMode é JSON, defina o tipo de campo JSON usando o campo type dentro do parâmetro column. Durante a gravação, o sistema tenta converter o valor da coluna de origem para o tipo especificado. Falhas na conversão resultam em dados incorretos (dirty data).

Valor válido

Descrição

JSON_STRING

Converte o valor da coluna de origem para string e grava no campo JSON. Exemplo: se o valor for o inteiro 123 e column for [{"name":"col1","type":"JSON_STRING"}], o valor gravado será {"col1":"123"}.

JSON_NUMBER

Converte o valor da coluna para número e grava no campo JSON. Exemplo: se o valor for a string 1.23 e column for [{"name":"col1","type":"JSON_NUMBER"}], o resultado será {"col1":1.23}.

JSON_BOOL

Converte o valor para Boolean e grava no campo JSON. Exemplo: se o valor for a string true e column for [{"name":"col1","type":"JSON_BOOL"}, o valor gravado será {"col1":true}

JSON_ARRAY

Converte o valor para um array JSON e grava no campo. Exemplo: se o valor for a string [1,2,3] e column for [{"name":"col1","type":"JSON_ARRAY"}], o resultado será {"col1":[1,2,3]}.

JSON_MAP

Converte o valor para um objeto JSON e grava no campo. Exemplo: se o valor for a string {"k1":"v1"} e column for [{"name":"col1","type":"JSON_MAP"}], o valor gravado será {"col1":{"k1":"v1"}}.

JSON_BASE64

Converte um array de bytes da coluna de origem para uma string codificada em BASE64 e grava no campo JSON. Exemplo: se o valor for um array de 2 bytes representado em hexadecimal como 0x01 0x02, e column for [{"name":"col1","type":"JSON_BASE64"}], o valor gravado será {"col1":"AQI="}.

JSON_HEX

Converte um array de bytes da coluna de origem para uma string hexadecimal e grava no campo JSON. Exemplo: se o valor for um array de 2 bytes representado em hexadecimal como 0x01 0x02, e column for [{"name":"col1","type":"JSON_HEX"}], o valor gravado será {"col1":"0102"}.