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.
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%.
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}ImportanteEssa 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.
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 |
|
||
|
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
Para mais detalhes sobre o procedimento, consulte Configure uma tarefa de sincronização offline na interface visual e Configure uma tarefa de sincronização offline no editor de código.
Para a lista completa de parâmetros e um exemplo de script para o editor de código, veja Apêndice: Exemplos de scripts e descrições de parâmetros.
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.
Se o cluster Kafka for uma instância do Alibaba Cloud Kafka, consulte Instruções de atualização do algoritmo de certificado SSL para baixar o arquivo de certificado truststore correto. A senha do truststore é KafkaOnsClient.
Caso o cluster Kafka seja uma instância EMR, acesse Usar criptografia SSL para conexões Kafka para obter o arquivo de certificado truststore adequado e a respectiva senha.
Para clusters auto-gerenciados, envie o certificado truststore correto e informe a senha apropriada.
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.
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.
|
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:
|
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
|
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.
|
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.
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: 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:
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:
|
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.
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 |
|
JSON_NUMBER |
Converte o valor da coluna para número e grava no campo JSON. Exemplo: se o valor for a string |
|
JSON_BOOL |
Converte o valor para Boolean e grava no campo JSON. Exemplo: se o valor for a string |
|
JSON_ARRAY |
Converte o valor para um array JSON e grava no campo. Exemplo: se o valor for a string |
|
JSON_MAP |
Converte o valor para um objeto JSON e grava no campo. Exemplo: se o valor for a string |
|
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 |
|
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 |