Todos os produtos
Search
Central de documentação

DataWorks:Fonte de dados Kafka

Última atualização: Aug 26, 2026

A fonte de dados Kafka oferece um canal bidirecional para leitura e gravação de dados no Kafka. Este tópico descreve os recursos 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 em versões do Kafka anteriores à 0.10.2, pois essas versões não permitem a recuperação de offsets de partição e suas estruturas de dados podem não suportar timestamps.

Leitura em tempo real

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

    Estime 1 CU por topic. Estime também os recursos com base no tráfego:

    • Para dados do Kafka sem compactaçã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.

    • Se os dados compactados do Kafka exigirem 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%.

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

Nota

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

Limitações

A fonte de dados Kafka suporta Serverless resource groups (recommended) e old versions of exclusive resource groups for Data Integration.

Leitura offline de uma única tabela

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 uma única tabela

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 um banco de dados inteiro

  • Tarefas de sincronização de dados em tempo real suportam Serverless resource groups (recommended) e old versions of exclusive resource groups for Data Integration.

  • Quando a tabela de origem possui uma 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 topic 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 são chave primária será usada como chave do registro Kafka.

  • Para assegurar que alterações da mesma chave primária sejam gravadas na mesma partição do Kafka em ordem, mesmo que o cluster Kafka retorne 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. É necessário 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 Appendix: Message format.

Tipos de campo suportados

O Kafka fornece armazenamento de dados não estruturados. Um registro Kafka geralmente inclui os seguintes campos: key, value, offset, timestamp, headers e partition. Quando o DataWorks lê ou grava dados no Kafka, ele processa os dados da seguinte forma.

Leitura de dados

Ao ler dados do Kafka, o DataWorks pode analisar os dados no formato JSON. A tabela a seguir 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

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

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.

  • Quando uma tarefa de sincronização em tempo real grava dados no Kafka, ela utiliza 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 Appendix: Message format.

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 codificada em UTF-8 "true" ou "false"

Time/Date

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

Numeric

String numérica codificada em UTF-8

Fluxo de bytes

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

text

String

String codificada em UTF-8

Boolean

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

Time/Date

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

Numeric

String numérica codificada em UTF-8

Fluxo de bytes

O fluxo de bytes é tratado como uma string codificada em UTF-8 e convertido em 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 do JSON

Time/Date

  • Para valores de tempo com precisão inferior a milissegundos: Convertido em um número inteiro JSON de 13 dígitos que representa o timestamp em milissegundos.

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

Numeric

Tipo numérico do JSON

Fluxo de bytes

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

text

String

String codificada em UTF-8

Boolean

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

Time/Date

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

Numeric

String numérica codificada em UTF-8

Fluxo de bytes

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

Sincronização em tempo real: Sincronização em tempo real de um 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 do JSON

Time/Date

Timestamp de 13 dígitos em milissegundos

Numeric

Valor numérico JSON

Fluxo de bytes

O fluxo de bytes é codificado em Base64 e depois convertido em 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 do JSON

Time/Date

Timestamp de 13 dígitos em milissegundos

Numeric

Valor numérico JSON

Fluxo de bytes

O fluxo de bytes é codificado em Base64 e depois convertido em 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 Data source configuration. 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.

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

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

Para obter instruções, consulte Configure a real-time synchronization task for a single table e Configure a real-time synchronization task for an entire database.

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. É necessário fazer upload do arquivo de certificado truststore do cliente e inserir a senha do truststore.

  • Se o cluster Kafka for uma instância do Alibaba Cloud Kafka, consulte SSL certificate algorithm upgrade instructions para baixar o arquivo de certificado truststore correto. A senha do truststore é KafkaOnsClient.

  • Caso o cluster Kafka seja uma instância EMR, consulte Use SSL encryption for Kafka connections para baixar o arquivo de certificado truststore correto e obter a senha do truststore.

  • Para um cluster auto-gerenciado, faça upload do certificado truststore correto e insira a senha adequada.

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 usa essas informações 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 Use SSL encryption for Kafka connections.

GSSAPI

Se você definir o Sasl Mechanism como GSSAPI ao configurar uma fonte de dados Kafka, será necessário fazer upload de 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 descrevem esses arquivos e as configurações necessárias de DNS e HOST.

Nota

Para um Serverless resource group, configure as informações de endereço do host usando resolução de DNS interno. Para mais informações, consulte Internal DNS resolution (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, a classe do componente de logon é fixa. Cada item de configuração subsequente é escrito no 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 a seguir mostra um formato típico de arquivo de configuração JAAS. Substitua os espaços reservados 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

    Deve ser definido como true.

    keyTab

    É possível especificar qualquer caminho. Quando a tarefa de sincronização é executada, o sistema baixa automaticamente o arquivo keytab enviado durante a configuração da fonte de dados para um caminho local. Em seguida, o sistema usa esse caminho local para o item de configuração keytab.

    storeKey

    Especifica se o cliente salva a chave. Pode ser definido como true ou false. Não afeta a sincronização de dados.

    serviceName

    Corresponde ao item de configuração sasl.kerberos.service.name no arquivo de configuração server.properties do servidor Kafka. Configure este item conforme necessário.

    principal

    O 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 de configuração no módulo é escrito no formato key=value.

    • O módulo [realms] especifica o endereço do Key Distribution Center (KDC). Ele pode conter vários submódulos realm. Cada submódulo realm começa com o nome do realm seguido por um sinal de igual (=).

    Em seguida, segue-se um conjunto de itens de configuração entre chaves. Cada item de configuração também é escrito no formato key=value. O código a seguir mostra um formato típico de arquivo de configuração Kerberos. Substitua os espaços reservados xxx pelas suas informações reais.

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

    Item de configuração

    Descrição

    [libdefaults].default_realm

    O realm padrão usado ao acessar nós do cluster Kafka. Geralmente é o mesmo realm do principal do cliente especificado no arquivo de configuração JAAS.

    Outros parâmetros de [libdefaults]

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

    [realms].nome do realm

    Deve ser igual ao realm do principal do cliente especificado no arquivo de configuração JAAS e ao [libdefaults].default_realm. Se o realm do principal do cliente no arquivo de configuração JAAS for diferente de [libdefaults].default_realm, será necessário incluir dois submódulos realms. Esses submódulos devem corresponder ao realm do principal do cliente no arquivo de configuração JAAS e ao [libdefaults].default_realm, respectivamente.

    [realms].nome do realm.kdc

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

  • Arquivo Keytab

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

    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 um grupo de recursos exclusivo

    Quando um cluster Kafka usa autenticação Kerberos, o KDC registra o principal de cada nó usando o hostname do nó. Quando um cliente se conecta a um nó do cluster Kafka, ele usa as configurações locais de DNS e HOST para derivar o principal do nó e, em seguida, solicita uma credencial de acesso para o nó ao KDC. Ao usar um grupo de recursos exclusivo para acessar um cluster Kafka com autenticação Kerberos habilitada, configure corretamente as definições de DNS e HOST para garantir que as credenciais de acesso aos nós do cluster possam ser obtidas do KDC:

    • Configurações de DNS

      Se você utilizar uma instância PrivateZone para resolução de nomes de domínio dos nós do cluster Kafka na VPC à qual o grupo de recursos exclusivo está anexado, adicione uma rota personalizada para os endereços IP 100.100.2.136 e 100.100.2.138 ao anexo da VPC. Isso garante que as configurações de resolução de nomes de domínio do PrivateZone para os nós do cluster Kafka se apliquem ao grupo de recursos exclusivo. No painel de navegação à esquerda do console do 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 resolução de nomes de domínio dos nós do cluster Kafka na VPC à qual o grupo de recursos exclusivo está anexado, adicione os mapeamentos de endereço IP para nome de domínio de cada nó do cluster Kafka na configuração de host. No painel de navegação à esquerda do console do 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 adicionar 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 você definir o Sasl Mechanism como PLAIN, 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, a classe do componente de logon é fixa. Cada item de configuração subsequente é escrito no 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. Um ponto e vírgula também deve ser adicionado após 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 a seguir mostra um formato típico de arquivo de configuração JAAS. Substitua os espaços reservados xxx pelas suas informações 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

O nome de usuário. Configure este item conforme necessário.

password

A senha. Configure este item conforme necessário.

Perguntas frequentes

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

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

Para configurar uma tarefa de sincronização em lote usando o editor de código, defina os parâmetros relacionados no script de acordo com os requisitos unificados de formato de script. Para mais informações, consulte Script mode configuration. As informações a seguir descrevem os parâmetros que devem ser configurados para fontes de dados ao utilizar o editor de código para tarefas de sincronização em lote.

Exemplo de script do Reader

A configuração JSON a seguir 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

O nome da fonte de dados. O editor de código suporta a adição de fontes de dados. O valor deste parâmetro deve ser idêntico ao nome da fonte de dados adicionada.

Sim

server

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

É possível configurar apenas um server, mas é necessário garantir que o DataWorks possa se conectar aos endereços IP de todos os brokers no cluster Kafka.

Sim

topic

O topic do Kafka. Um topic é uma agregação de fluxos de mensagens processados pelo Kafka.

Sim

column

Os dados do Kafka a serem lidos. Colunas constantes, colunas de dados e colunas de atributos são suportadas.

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

  • Coluna de dados

    • Se seus dados estiverem no formato JSON, é possível obter as propriedades do objeto JSON, como ["event_id"].

    • Se seus dados estiverem no formato JSON, é possível obter as subpropriedades aninhadas do objeto JSON, como ["tag.desc"].

  • Coluna de atributo

    • __key__: A chave da mensagem.

    • __value__: O conteúdo completo da mensagem.

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

    • __headers__: Os cabeçalhos da mensagem atual.

    • __offset__: O offset da mensagem atual.

    • __timestamp__: O timestamp da mensagem atual.

    O código a seguir fornece um exemplo completo.

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

Sim

keyType

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

Não

valueType

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

Não

beginDateTime

A hora inicial para consumo de dados. Este parâmetro especifica o limite esquerdo do intervalo de tempo, que é inclusivo. É uma string de tempo no formato yyyymmddhhmmss. Use este parâmetro com scheduling parameters. Para mais informações, consulte Supported formats of scheduling parameters.

Nota

Este recurso é suportado no Kafka 0.10.2 e versões posteriores.

Especifique este parâmetro ou beginOffset.

Nota

beginDateTime e endDateTime são usados em conjunto.

endDateTime

A hora final para consumo de dados. Este parâmetro especifica o limite direito do intervalo de tempo, que é exclusivo. É uma string de tempo no formato yyyymmddhhmmss. Use este parâmetro com scheduling parameters. Para mais informações, consulte Supported formats of scheduling parameters.

Nota

Este recurso é suportado no Kafka 0.10.2 e versões posteriores.

Especifique este parâmetro ou endOffset.

Nota

endDateTime e beginDateTime são usados em conjunto.

beginOffset

O offset inicial para consumo de dados. Configure-o das seguintes formas:

  • Um número, como 15553274, que indica o offset inicial para consumo.

  • seekToBeginning: indica que os dados são consumidos a partir do offset mais antigo.

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

  • seekToEnd: indica que os dados são consumidos a partir do offset mais recente. Isso resultará na leitura de dados vazios.

Especifique este parâmetro ou beginDateTime.

endOffset

O offset final para consumo de dados. Usado para controlar quando a tarefa de consumo de dados é encerrada.

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ções automáticas de offset, recomendamos o seguinte:

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

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

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

partition

Um topic do Kafka possui múltiplas partições (partition). Por padrão, uma tarefa de sincronização de dados lê dados de um intervalo de offsets que cobre todas as partições do topic. 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, especifique 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 é definido como STRING, a codificação especificada por este parâmetro é usada para analisar a string.

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

waitTIme

O tempo máximo, em segundos, que o objeto consumidor aguarda para extrair dados do Kafka em uma única tentativa.

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

stopWhenPollEmpty

Os valores válidos são true e false. Se este parâmetro for definido como true e o consumidor extrair dados vazios do Kafka (geralmente porque todos os dados no topic foram lidos, ou devido a problemas de rede ou disponibilidade do cluster Kafka), a tarefa para imediatamente. Caso contrário, ela tenta novamente até que os dados sejam lidos.

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

stopWhenReachEndOffset

Este parâmetro só entra em vigor quando stopWhenPollEmpty é true. Os valores válidos são true e false.

  • Se este parâmetro for definido como true e o consumidor não receber dados de uma solicitação poll, ele verifica se o offset mais recente nas partições do topic foi atingido. Se o offset mais recente tiver sido atingido para todas as partições, a tarefa para imediatamente. Caso contrário, continua tentando extrair dados do topic.

  • Se este parâmetro for definido como false e o consumidor não receber dados de uma solicitação poll, ele não realiza a verificação e interrompe a tarefa imediatamente.

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

Nota

Esta configuração fornece compatibilidade com versões anteriores. Versões do Kafka anteriores à 0.10.2 não suportam a 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

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

fetch.max.wait.ms

O tempo máximo, em milissegundos, que o broker aguarda até que os dados estejam disponíveis antes de responder a uma solicitação de busca. O valor padrão é 500. O broker responde quando a condição fetch.min.bytes ou fetch.max.wait.ms for atendida.

max.partition.fetch.bytes

Especifica o número máximo de bytes que o broker pode retornar ao consumidor de cada partition. O valor padrão é 1 MB.

session.timeout.ms

Especifica o tempo que o consumidor pode ficar desconectado do servidor antes de parar de receber serviços. O valor padrão é 30 segundos.

auto.offset.reset

A ação que o consumidor executa ao ler sem offset ou com um offset inválido (porque o consumidor esteve inativo por muito tempo e o registro com o offset expirou e foi excluído). O valor padrão é none, o que significa que o offset não é redefinido automaticamente. Altere para earliest para que o consumidor leia os registros da partition a partir do offset mais antigo.

max.poll.records

O número de mensagens que podem ser retornadas por uma única chamada ao método poll.

key.deserializer

O método de desserialização para a chave da mensagem, como org.apache.kafka.common.serialization.StringDeserializer.

value.deserializer

O método de desserialização para o valor dos dados, como org.apache.kafka.common.serialization.StringDeserializer.

ssl.truststore.location

O caminho do certificado raiz SSL.

ssl.truststore.password

A senha para o armazenamento de certificados raiz. Se estiver usando o Alibaba Cloud Kafka, defina como KafkaOnsClient.

security.protocol

O protocolo de acesso. Atualmente, apenas o protocolo SASL_SSL é suportado.

sasl.mechanism

O método de autenticação SASL. Se estiver usando o Alibaba Cloud Kafka, use PLAIN.

java.security.auth.login.config

O caminho do arquivo de autenticação SASL.

Exemplo de script do Writer

O código a seguir 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

O nome da fonte de dados. O editor de código suporta a adição de fontes de dados. O valor deste parâmetro deve ser idêntico ao nome da fonte de dados adicionada.

Sim

server

O endereço do servidor Kafka no formato IP:port.

Sim

topic

O topic do Kafka. É uma categoria para diferentes fluxos de mensagens processados pelo Kafka.

Cada mensagem publicada em um cluster Kafka possui uma categoria, chamada de topic. Um topic é uma coleção de um grupo de mensagens.

Sim

valueIndex

A coluna no writer do Kafka usada como valor. Se não for especificado, todas as colunas serão concatenadas para formar o valor por padrão. O separador é especificado por fieldDelimiter.

Não

writeMode

Quando valueIndex não está configurado, este parâmetro determina o formato para concatenar todas as colunas do registro de origem para formar o valor do registro Kafka. Os valores válidos são text e JSON. O valor padrão é text.

  • Se definido como text, todas as colunas são concatenadas usando o separador especificado por fieldDelimiter.

  • Se definido como JSON, todas as colunas são concatenadas em uma string JSON com base nos nomes de campo especificados pelo parâmetro column.

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

Se valueIndex estiver configurado, este parâmetro será inválido.

Não

column

Os campos na tabela de destino onde os dados serão gravados, separados por vírgulas. Por exemplo: "column": ["id", "name", "age"].

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

  • Se o número de colunas no registro de origem for maior que o número de nomes de campo configurados em column, os dados serão truncados durante a gravação. Por exemplo:

    Se um registro de origem tiver três colunas com valores a, b e c, e column estiver configurado como [{"name":"col1","type":"JSON_STRING"},{"name":"col2","type":"JSON_STRING"}], o valor do registro Kafka gravado será a string {"col1":"a","col2":"b"}.

  • Se o número de colunas no registro de origem for menor que o número de campos especificados em column, os campos extras serão preenchidos com null ou a string especificada por nullValueFormat. Por exemplo:

    Se um registro de origem tiver duas colunas com valores a e b, e column estiver configurado como [{"name":"col1","type":"JSON_STRING"},{"name":"col2","type":"JSON_STRING"},{"name":"col3","type":"JSON_STRING"}], o valor do registro Kafka gravado será a string {"col1":"a","col2":"b","col3":null}. Se valueIndex estiver configurado, ou se writeMode estiver definido como text, este parâmetro será inválido.

  • Se o tipo de campo JSON não estiver configurado, o tipo de campo padrão será JSON_STRING.

  • Para os valores válidos do tipo de campo JSON, consulte Appendix: JSON field types.

Se valueIndex estiver configurado, ou se writeMode estiver definido como text, este parâmetro será inválido.

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

partition

Especifica o número da partição no topic do Kafka onde os dados serão gravados. Deve ser um número inteiro maior ou igual a 0.

Não

keyIndex

A coluna no writer do Kafka usada como chave.

O valor do parâmetro keyIndex deve ser um número inteiro maior ou igual a 0. Caso contrário, a tarefa falhará.

Não

keyIndexes

Uma matriz dos números ordinais das colunas no registro de origem que são usadas como chave para o registro Kafka.

O número ordinal da coluna começa em 0. Por exemplo, [0,1,2] concatenará os valores de todos os números de coluna configurados com vírgulas para formar a chave do registro Kafka. Se não for especificado, a chave do registro Kafka será null, e os dados serão gravados nas partições do topic de forma round-robin. Especifique apenas este parâmetro ou keyIndex.

Não

fieldDelimiter

Quando writeMode está definido como text e valueIndex não está configurado, todas as colunas do registro de origem são concatenadas usando o separador de coluna especificado por este parâmetro para formar o valor do registro Kafka. Configure um único caractere ou múltiplos caracteres como separador. Caracteres Unicode podem ser configurados no formato \u0001. Caracteres de escape como \t e \n são suportados. O valor padrão é \t.

Se writeMode não estiver definido como text ou se valueIndex estiver configurado, este parâmetro será inválido.

Não

keyType

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

Sim

valueType

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

Sim

nullKeyFormat

Se o valor da coluna de origem especificada por keyIndex ou keyIndexes for null, ele será substituído pela string especificada neste parâmetro. Se não estiver configurado, nenhuma substituição será feita.

Não

nullValueFormat

Se o valor de uma coluna de origem for null, ele será substituído pela string especificada neste parâmetro ao montar o valor do registro Kafka. Se você não especificar este parâmetro, nenhuma substituição será feita.

Não

acks

A configuração acks ao inicializar o produtor Kafka. Determina o método de confirmação para gravações bem-sucedidas. Por padrão, o parâmetro acks é definido como all. Os valores válidos para acks são:

  • 0: Nenhuma confirmação para gravações bem-sucedidas.

  • 1: Confirmação para gravação bem-sucedida na réplica primária.

  • all: Confirmação para gravação bem-sucedida em todas as réplicas.

Não

Apêndice: Definição do 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 dados de origem são gravados em um topic do Kafka no formato JSON. Primeiro, todos os dados existentes na tabela de origem especificada são gravados no topic correspondente do Kafka. Em seguida, a tarefa inicia a sincronização em tempo real para gravar continuamente dados incrementais no topic. Informações de alterações DDL incrementais da tabela de origem também são gravadas no topic do Kafka no formato JSON. Obtenha o status e as informações de alteração das mensagens gravadas no Kafka. Para mais informações, consulte Appendix: Message format.

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 está definido como JSON, defina o tipo de campo JSON usando o campo type dentro do parâmetro column. Durante uma operação de gravação, o sistema tenta converter o valor da coluna do registro de origem para o tipo especificado. Uma falha na conversão de tipo resulta em dados incorretos.

Valor válido

Descrição

JSON_STRING

Converte o valor da coluna do registro de origem em uma string e o grava no campo JSON. Por exemplo, se o valor da coluna do registro de origem for o número inteiro 123 e column estiver configurado como [{"name":"col1","type":"JSON_STRING"}], o valor gravado no registro Kafka será a string {"col1":"123"}.

JSON_NUMBER

Converte o valor da coluna do registro de origem em um número e o grava no campo JSON. Por exemplo, se o valor da coluna do registro de origem for a string 1.23 e column estiver configurado como [{"name":"col1","type":"JSON_NUMBER"}], o valor gravado no registro Kafka será a string {"col1":1.23}.

JSON_BOOL

Converte o valor da coluna do registro de origem em um valor booleano e o grava no campo JSON. Por exemplo, se o valor da coluna do registro de origem for a string true e column estiver configurado como [{"name":"col1","type":"JSON_BOOL"}, o valor do registro Kafka gravado será a string {"col1":true}

JSON_ARRAY

Converte o valor da coluna do registro de origem em uma matriz JSON e o grava no campo JSON. Por exemplo, se o valor da coluna do registro de origem for a string [1,2,3] e column estiver configurado como [{"name":"col1","type":"JSON_ARRAY"}], o valor gravado no registro Kafka será a string {"col1":[1,2,3]}.

JSON_MAP

Converte o valor da coluna do registro de origem em um objeto JSON e o grava no campo JSON. Por exemplo, se o valor da coluna do registro de origem for a string {"k1":"v1"} e column estiver configurado como [{"name":"col1","type":"JSON_MAP"}], o valor gravado no registro Kafka será a string {"col1":{"k1":"v1"}}.

JSON_BASE64

Converte uma matriz de bytes da coluna de origem em uma string codificada em BASE64 e a grava no campo JSON. Por exemplo, se o valor da coluna do registro de origem for uma matriz de 2 bytes representada em hexadecimal como 0x01 0x02, e column estiver configurado como [{"name":"col1","type":"JSON_BASE64"}], o valor gravado no registro Kafka será a string {"col1":"AQI="}.

JSON_HEX

Converte uma matriz de bytes da coluna de origem em uma string hexadecimal e a grava no campo JSON. Por exemplo, se o valor da coluna do registro de origem for uma matriz de 2 bytes representada em hexadecimal como 0x01 0x02, e column estiver configurado como [{"name":"col1","type":"JSON_HEX"}], o valor gravado no registro Kafka será a string {"col1":"0102"}.