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.
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%.
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}ImportanteEssa 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.
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 |
|
||
|
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
Para mais detalhes sobre o procedimento, consulte Configure an offline synchronization task in the codeless UI e Configure an offline synchronization task in the code editor.
Para a lista completa de parâmetros e um exemplo de script para o editor de código, consulte Appendix: Script demos and parameter descriptions.
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.
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.
|
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:
|
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
|
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.
|
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.
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: 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,
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:
|
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.
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 |
|
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 |
|
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 |
|
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 |
|
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 |
|
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 |
|
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 |