Saiba como usar o conector do Log Service (SLS).
Contexto
O Simple Log Service é um serviço completo para dados de log. Ele permite coletar, consumir, enviar, consultar e analisar dados de log com eficiência, melhorando a eficiência de O&M e possibilitando o processamento de grandes volumes de dados.
A tabela a seguir lista as capacidades do conector SLS.
|
Categoria |
Descrição |
|
Tipos suportados |
Tabela de origem e tabela de destino |
|
Modo de execução |
Apenas modo streaming |
|
Métricas específicas do conector |
N/A |
|
Formato de dados |
N/A |
|
Tipo de API |
SQL, DataStream API e API yaml de ingestão de dados |
|
Atualização ou exclusão de dados na tabela de destino |
As tabelas de destino aceitam apenas adições; não é possível atualizar ou excluir dados. |
Recursos
O conector de source SLS lê diretamente os campos de atributos das mensagens. A tabela a seguir lista os campos suportados.
|
Parâmetro |
Tipo |
Descrição |
|
__source__ |
STRING METADATA VIRTUAL |
Origem da mensagem. |
|
__topic__ |
STRING METADATA VIRTUAL |
Tópico da mensagem. |
|
__timestamp__ |
BIGINT METADATA VIRTUAL |
Hora do log. |
|
__tag__ |
MAP<VARCHAR, VARCHAR> METADATA VIRTUAL |
Tag da mensagem. Por exemplo, para o atributo |
Pré-requisitos
Certifique-se de ter criado um Project do Log Service e um Logstore. Para mais informações, consulte Criar um Project e um Logstore.
Limitações
Somente o Ververica Runtime (VVR) 11,1 e versões posteriores suportam o uso do SLS como fonte síncrona para ingestão de dados definida em yaml.
O conector SLS garante apenas a semântica at-least-once.
Não defina o
source parallelismmaior que o número de shards, pois isso desperdiça recursos. Além disso, no Ververica Runtime (VVR) 8.0.5 e anteriores, uma alteração no número de shards pode causar falha no recurso automático defailover, impedindo o consumo de alguns shards.
SQL
Sintaxe
CREATE TABLE sls_table(
a INT,
b INT,
c VARCHAR
) WITH (
'connector' = 'sls',
'endPoint' = '<yourEndPoint>',
'project' = '<yourProjectName>',
'logStore' = '<yourLogStoreName>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}'
);
Opções With
-
Gerais
Parâmetro
Descrição
Tipo
Obrigatório
Padrão
Observações
connector
O conector a ser usado.
String
Sim
Nenhum
Defina como
sls.endPoint
O endpoint do Log Service (SLS).
String
Sim
Nenhum
Especifique o endereço de acesso vpc do Log Service (SLS). Para mais informações, consulte Endpoints de serviço.
Nota-
Por padrão, o Realtime Compute for Apache Flink não acessa a internet. Para habilitar o acesso à internet a partir da sua Virtual Private Cloud (vpc), use um NAT Gateway. Para mais informações, consulte Como acesso a Internet?.
-
Recomendamos não acessar o SLS pela internet. Caso seja necessário, use HTTPS e ative a aceleração de transferência. Para mais informações, consulte Gerenciar aceleração de transferência.
project
O nome do projeto SLS.
String
Sim
Nenhum
N/A
logStore
O nome do Logstore ou MetricStore.
String
Sim
Nenhum
Os dados em um Logstore são consumidos da mesma forma que os dados em um MetricStore.
accessId
O AccessKey ID da sua conta Alibaba Cloud.
String
Sim
Nenhum
Para mais informações, consulte Obter um par de AccessKey.
ImportantePara evitar a exposição do seu par de AccessKey, recomendamos o uso de variáveis para especificar o AccessKey ID e o AccessKey secret. Para mais informações, consulte Variáveis de projeto.
accessKey
O AccessKey secret da sua conta Alibaba Cloud.
String
Sim
Nenhum
-
-
Específicas da origem
Parâmetro
Descrição
Tipo
Obrigatório
Padrão
Observações
enableNewSource
Define se deve ser usada a nova fonte de dados que implementa a interface FLIP-27.
Boolean
Não
false
A nova fonte adapta-se automaticamente a alterações nos shards e distribui os shards da forma mais uniforme possível entre todas as subtarefas de origem.
Importante-
Esta opção é suportada apenas no VVR 8.0.9 e posteriores. O valor padrão é
truepara VVR 11,1 e posteriores. -
Se você alterar o valor desta opção, não será possível restaurar jobs a partir de um estado salvo. Para retomar o consumo a partir de um offset histórico, inicie primeiro o job com a opção
consumerGrouppara registrar o progresso de consumo no grupo de consumidores do SLS. Em seguida, defina a opçãoconsumeFromCheckpointcomotruee reinicie o job sem estado. -
Se um Logstore possuir shards somente leitura, algumas subtarefas podem continuar solicitando dados de outros shards após concluírem os seus próprios. Isso pode causar desequilíbrio na carga de trabalho, afetando o desempenho. Para mitigar esse problema, ajuste o paralelismo, otimize a estratégia de agendamento ou mescle shards pequenos para reduzir o número total de shards e simplificar a alocação de tarefas.
shardDiscoveryIntervalMs
O intervalo para descoberta dinâmica de shards.
Long
Não
60000
Defina esta opção com um valor negativo para desativar a descoberta dinâmica. Unidade: milissegundos.
Nota-
O valor deve ser maior ou igual a 60.000 milissegundos (1 minuto).
-
Esta opção só tem efeito quando
enableNewSourceestá definido comotrue. -
Esta opção é suportada apenas no VVR 8.0.9 e posteriores.
startupMode
O modo de inicialização da tabela de origem.
String
Não
timestamp
-
timestamp(padrão): Consome logs a partir da hora de início especificada. -
latest: Inicia o consumo de logs a partir do offset mais recente. -
earliest: Inicia o consumo de logs a partir do offset mais antigo. -
consumer_group: Inicia o consumo de logs a partir do offset registrado pelo grupo de consumidores. Se o grupo de consumidores não tiver registrado um offset de consumo para um shard, o consumo começa a partir do offset mais antigo.
Importante-
Para versões do VVR anteriores à 11,1, o valor consumer_group não é suportado. É necessário definir
consumeFromCheckpointcomotrue. Nesse caso, o consumo de logs inicia a partir do offset registrado pelo grupo de consumidores especificado, e a configuração do modo de inicialização não terá efeito.
startTime
A hora de início para o consumo de logs.
String
Não
Hora atual
O formato é
yyyy-MM-dd hh:mm:ss.Isso só tem efeito quando
startupModeestá definido comotimestamp.NotaAs opções
startTimeestopTimebaseiam-se no atributo__receive_time__no SLS, e não no atributo__timestamp__.stopTime
A hora de término para o consumo de logs.
String
Não
Nenhum
O formato é
yyyy-MM-dd hh:mm:ss.Nota-
Esta opção serve apenas para consumir logs históricos e deve ser definida com uma hora no passado. Se definida para uma hora futura, o consumo poderá parar prematuramente caso nenhum novo log seja gravado, resultando em interrupção do fluxo de dados sem mensagens de erro.
-
Se desejar que o job Flink encerre após o consumo de todos os logs, também defina
exitAfterFinishcomotrue.
consumerGroup
O nome do grupo de consumidores.
String
Não
Nenhum
Um grupo de consumidores registra o progresso de consumo. É possível especificar um nome personalizado sem formato fixo.
NotaJobs Flink diferentes devem usar grupos de consumidores distintos. Se múltiplos jobs Flink usarem o mesmo grupo de consumidores, eles não se coordenarão e cada job consumirá todos os dados. Isso ocorre porque o Flink não usa o grupo de consumidores do SLS para atribuição de partições ao consumir dados do SLS. Como resultado, cada consumidor processa mensagens independentemente, mesmo compartilhando o mesmo grupo de consumidores.
consumeFromCheckpoint
Define se o consumo deve iniciar a partir de um checkpoint do grupo de consumidores.
String
Não
false
-
true: Também é necessário especificar um grupo de consumidores. O programa Flink inicia o consumo de logs a partir do checkpoint salvo no grupo de consumidores. Se o grupo não possuir um checkpoint correspondente, o consumo começa a partir do valor configurado em startTime. -
false(valor padrão): Não inicia o consumo de logs a partir do checkpoint salvo para o grupo de consumidores especificado.
ImportanteEste parâmetro não é mais suportado no VVR 11,1 e posteriores. Para essas versões, defina a opção
startupModecomoconsumer_group.maxRetries
Número de tentativas após falha na leitura do SLS.
String
Não
3
N/A
batchGetSize
Quantidade de grupos de logs lidos em uma única solicitação.
String
Não
100
A configuração
batchGetSizenão pode exceder 1000; caso contrário, um erro será reportado.exitAfterFinish
Define se o job Flink encerra após o consumo de todos os dados.
String
Não
false
-
true: O programa Flink encerra após o consumo de todos os dados. -
false(padrão): O programa Flink não encerra após a conclusão do consumo de dados.
query
ImportanteEsta opção foi descontinuada no VVR 11.3, mas versões posteriores permanecem compatíveis.
Instrução de consulta para pré-processamento de dados antes do consumo.
String
Não
Nenhum
Use esta opção para filtrar dados do SLS antes que o Flink os consuma. Isso reduz custos e aumenta a velocidade de processamento.
Por exemplo,
'query' = '*| where request_method = ''GET'''indica que, antes de ler dados do SLS, o Flink seleciona apenas os registros onde o valor do camporequest_methodé 'GET'.NotaEsta opção usa a linguagem SPL do Log Service (SLS). Para mais informações, consulte Sintaxe SPL.
Importante-
Esta opção é suportada apenas no VVR 8.0.1 e posteriores.
-
Este recurso gera cobranças do Log Service (SLS). Para mais informações, consulte Faturamento do Log Service.
processor
Nome do processador SLS para pré-processamento de dados. Se tanto
queryquantoprocessorforem especificados,queryterá precedência eprocessorserá ignorado.String
Não
Nenhum
Esta opção filtra dados do SLS antes do consumo pelo Flink, o que reduz custos e melhora a velocidade de processamento. Recomendamos o uso de
processorem vez dequery.Por exemplo,
'processor' = 'test-filter-processor'indica que o processador SLS filtra os dados antes que o Flink leia do SLS.NotaEsta opção usa a linguagem SPL do Log Service (SLS). Para mais informações, consulte Sintaxe SPL. Para saber como criar ou atualizar um processador SLS, consulte Gerenciar processadores.
ImportanteEsta opção é suportada apenas no VVR 11.3 e posteriores.
Este recurso gera cobranças do Log Service (SLS). Para mais informações, consulte Faturamento do Log Service.
-
-
Específicas do destino
Parâmetro
Descrição
Tipo
Obrigatório
Padrão
Observações
topicField
Especifica um campo cujo valor sobrescreve o atributo
__topic__, que indica o tópico do log.String
Não
Nenhum
O valor desta opção deve ser um campo existente na tabela.
timeField
Especifica um campo cujo valor sobrescreve o atributo
__timestamp__, que indica a hora de gravação do log.String
Não
Hora atual
O valor desta opção deve ser um campo
INTexistente na tabela. Se esta opção não for especificada, a hora atual será usada.sourceField
Especifica um campo cujo valor sobrescreve o atributo
__source__, que indica a origem do log, como o endereço IP da máquina que gerou o log.String
Não
Nenhum
O valor desta opção deve ser um campo existente na tabela.
partitionField
Especifica um campo para particionamento. Um hash do valor deste campo determina qual shard recebe os dados, garantindo que registros com o mesmo hash sejam enviados ao mesmo shard.
String
Não
Nenhum
Se esta opção não for especificada, cada registro será gravado aleatoriamente em um shard disponível.
buckets
Se
partitionFieldfor especificado, esta opção define o número de buckets para mapeamento dos valores de hash.String
Não
64
O valor deve ser uma potência de 2 no intervalo [1, 256]. O número de buckets deve ser maior ou igual ao número de shards. Caso contrário, alguns shards podem não receber dados.
flushIntervalMs
O intervalo de gravação de dados.
String
Não
2000
Unidade: milissegundos.
writeNullProperties
Define se valores nulos devem ser gravados como strings vazias no SLS.
Boolean
Não
true
-
true(valor padrão): Grava valores nulos no log como strings vazias. -
false: Campos avaliados como nulo não são gravados no log.
NotaEsta opção é suportada apenas no VVR 8.0.6 e posteriores.
-
Mapeamentos de tipo
|
Tipo Flink |
Tipo SLS |
|
BOOLEAN |
STRING |
|
VARBINARY |
|
|
VARCHAR |
|
|
TINYINT |
|
|
INTEGER |
|
|
BIGINT |
|
|
FLOAT |
|
|
DOUBLE |
|
|
DECIMAL |
Ingestão de dados (Beta)
Limitações
Este recurso é suportado apenas pelas versões 11,1 e posteriores do Realtime Compute for Apache Flink.
Sintaxe
source:
type: sls
name: SLS Source
endpoint: <endpoint>
project: <project>
logstore: <logstore>
accessId: <accessId>
accessKey: <accessKey>
Parâmetros
|
Parâmetro |
Descrição |
Tipo |
Obrigatório |
Padrão |
Observações |
||||||||||
|
type |
O tipo da fonte de dados. |
String |
Sim |
Nenhum |
O valor deve ser |
||||||||||
|
endpoint |
O endpoint do Log Service (SLS). |
String |
Sim |
Nenhum |
O endereço de acesso vpc do Log Service (SLS). Para mais informações, consulte Endpoints de Serviço. Nota
|
||||||||||
|
accessId |
O AccessKey ID da sua conta Alibaba Cloud. |
String |
Sim |
Nenhum |
Para mais informações, consulte Como visualizo as informações de AccessKey ID e AccessKey secret?. Importante
Para evitar a exposição das suas informações de AccessKey, recomendamos o uso de uma variável de projeto para especificar o valor do AccessKey. Para mais informações, consulte Variáveis de projeto. |
||||||||||
|
accessKey |
O AccessKey secret da sua conta Alibaba Cloud. |
String |
Sim |
Nenhum |
|||||||||||
|
project |
O nome do projeto do Log Service (SLS). |
String |
Sim |
Nenhum |
Nenhum |
||||||||||
|
logStore |
O nome do Logstore ou Metricstore do SLS. |
String |
Sim |
Nenhum |
Os dados em um Logstore são consumidos da mesma forma que os dados em um Metricstore. |
||||||||||
|
schema.inference.strategy |
A estratégia de inferência de schema. |
String |
Não |
continuous |
|
||||||||||
|
maxPreFetchLogGroups |
Número máximo de grupos de logs a serem lidos de cada shard para a inferência inicial de schema. |
Integer |
Não |
50 |
Antes de ler e processar dados, o conector pré-consome um número especificado de grupos de logs de cada shard para inicializar as informações de schema. |
||||||||||
|
shardDiscoveryIntervalMs |
Intervalo, em milissegundos, para descoberta dinâmica de alterações nos shards. |
Long |
Não |
60000 |
Defina este parâmetro com um valor negativo para desativar a descoberta dinâmica. Nota
O valor deve ser maior ou igual a 60.000 milissegundos (1 minuto). |
||||||||||
|
startupMode |
O modo de inicialização. |
String |
Não |
Nenhum |
|
||||||||||
|
startTime |
A hora de início para o consumo de logs. |
String |
Não |
Hora atual |
O formato é Este parâmetro só tem efeito quando Nota
Os parâmetros |
||||||||||
|
stopTime |
A hora de término para o consumo de logs. |
String |
Não |
Nenhum |
O formato é Nota
Se desejar que o job Flink encerre após o consumo de todos os logs, também defina |
||||||||||
|
consumerGroup |
O nome do grupo de consumidores. |
String |
Não |
Nenhum |
Um grupo de consumidores registra o progresso de consumo. É possível especificar qualquer nome personalizado. |
||||||||||
|
batchGetSize |
Número de grupos de logs lidos por solicitação. |
Integer |
Não |
100 |
O valor de batchGetSize não pode exceder 1.000. Caso contrário, ocorrerá um erro. |
||||||||||
|
maxRetries |
Número de tentativas caso a leitura do Log Service (SLS) falhe. |
Integer |
Não |
3 |
Nenhum |
||||||||||
|
exitAfterFinish |
Define se o job Flink encerra após o consumo de todos os dados. |
Boolean |
Não |
false |
|
||||||||||
|
query |
Instrução de pré-processamento para consumo de dados do Log Service (SLS). |
String |
Não |
Nenhum |
Use este parâmetro para filtrar dados no Log Service (SLS) antes do consumo, visando economizar custos e aumentar a velocidade de processamento. Por exemplo, Nota
A consulta deve usar a sintaxe SPL do Log Service. Para mais informações, consulte Sintaxe SPL. Importante
|
||||||||||
|
compressType |
O tipo de compressão para o Log Service (SLS). |
String |
Não |
Nenhum |
Os tipos de compressão suportados incluem:
|
||||||||||
|
timeZone |
O fuso horário para |
String |
Não |
Nenhum |
Por padrão, nenhum deslocamento é adicionado. |
||||||||||
|
regionId |
A região onde o Log Service (SLS) está implantado. |
String |
Não |
Nenhum |
Para mais informações, consulte Regiões suportadas. |
||||||||||
|
signVersion |
A versão da assinatura de solicitação para o Log Service (SLS). |
String |
Não |
Nenhum |
Para mais informações, consulte Assinaturas de solicitação. |
||||||||||
|
shardModDivisor |
O divisor usado na leitura de shards do Logstore. |
Int |
Não |
-1 |
Para mais informações, consulte Shards. |
||||||||||
|
shardModRemainder |
O resto usado na leitura de shards do Logstore. |
Int |
Não |
-1 |
Para mais informações, consulte Shards. |
||||||||||
|
metadata.list |
As colunas de metadados a serem passadas para jobs downstream. |
String |
Não |
Nenhum |
Os campos de metadados disponíveis incluem |
||||||||||
|
decode.table-id.fields |
Especifica campos cujos valores são usados para gerar um Table ID ao analisar dados de log do Log Service (SLS). |
String |
Não |
Nenhum |
Múltiplos campos são separados por vírgula
Nota
Este parâmetro é suportado nas versões 11.6 e posteriores do Realtime Compute for Apache Flink. |
||||||||||
|
fixed-types |
Especifica os tipos de dados para campos específicos ao analisar dados de log do Log Service (SLS). |
String |
Não |
Nenhum |
Ao analisar dados, especifique os tipos para campos específicos. Use vírgula Nota
Este parâmetro é suportado nas versões 11.6 e posteriores do Realtime Compute for Apache Flink. |
||||||||||
|
timestamp-format.standard |
O formato para campos de timestamp nos dados de log do Log Service (SLS). |
String |
Não |
SQL |
Valores válidos:
Nota
Este parâmetro é suportado nas versões 11.6 e posteriores do Realtime Compute for Apache Flink. |
||||||||||
|
ingestion.ignore-errors |
Define se erros ocorridos durante a análise de dados devem ser ignorados. |
Boolean |
Não |
false |
Nota
Este parâmetro é suportado nas versões 11.6 e posteriores do Realtime Compute for Apache Flink. |
||||||||||
|
ingestion.error-tolerance.max-count |
Se |
Integer |
Não |
-1 |
Este parâmetro só tem efeito quando Nota
Este parâmetro é suportado nas versões 11.6 e posteriores do Realtime Compute for Apache Flink. |
Reutilizar um catálogo existente
A partir da versão 11,5 do Realtime Compute for Apache Flink, é possível referenciar um Catálogo SLS integrado criado na página Data Management em jobs de ingestão de dados Flink CDC. Isso evita a necessidade de especificar manualmente as propriedades de conexão.
source:
type: sls
using.built-in-catalog: sls_catalog
Atualmente, os jobs de ingestão de dados podem reutilizar automaticamente os seguintes parâmetros do Catálogo SLS integrado:
endpoint
project
accessId
accessKey
Para substituir esses parâmetros reutilizados automaticamente, defina-os explicitamente na configuração yaml. Os parâmetros definidos no arquivo yaml têm precedência.
Mapeamento de tipos de dados
Se fixed-types não estiver configurado, aplica-se o seguinte mapeamento de tipos de dados:
|
Tipo SLS |
Tipo CDC |
|
STRING |
STRING |
Se fixed-types estiver configurado, o sistema analisa os dados usando os tipos especificados.
Inferência e evolução de schema
-
Pré-consumo de dados de shard e inicialização de schema
O conector SLS mantém o schema do Logstore que está lendo. Antes de ler dados do Logstore, o conector pré-consome até
maxPreFetchLogGroupsgrupos de logs de cada shard. Ele analisa o schema de cada entrada de log e os mescla para inicializar o schema da tabela. Um evento de criação de tabela é então gerado com base neste schema inicial antes do início do consumo de dados.NotaPara cada shard, o conector tenta consumir dados começando de uma hora antes da hora atual para analisar o schema do log.
-
Informações de chave primária
Os logs do Log Service (SLS) não contêm informações de chave primária. Adicione manualmente uma chave primária à tabela usando regras de transformação:
transform: - source-table: <project>.<logstore> projection: * primary-keys: key1, key2 -
Inferência de schema e alterações de schema
Após a inicialização do schema, se schema.inference.strategy estiver definido como
static, o conector SLS analisa cada entrada de log com base no schema inicial e não gera eventos de alteração de schema. Se schema.inference.strategy estiver definido comocontinuous, o conector analisa cada entrada de log, infere as colunas físicas e as compara com o schema atual. Se o schema inferido for inconsistente com o schema atual, os schemas são mesclados de acordo com as seguintes regras:Se as colunas físicas inferidas contiverem campos que não estão no schema atual, o conector adiciona esses campos ao schema e gera um evento para adicionar colunas anuláveis.
Se as colunas físicas inferidas não possuírem campos existentes no schema atual, o conector retém esses campos, preenche seus dados com NULL e não gera um evento de exclusão de coluna.
O conector SLS infere o tipo de dados de todos os campos em cada entrada de log como String. Atualmente, apenas a adição de novas colunas é suportada. O conector anexa novas colunas ao final do schema e as define como anuláveis.
Exemplos de código
-
SQL para tabelas de origem e destino
CREATE TEMPORARY TABLE sls_input( `time` BIGINT, url STRING, dt STRING, float_field FLOAT, double_field DOUBLE, boolean_field BOOLEAN, `__topic__` STRING METADATA VIRTUAL, `__source__` STRING METADATA VIRTUAL, `__timestamp__` STRING METADATA VIRTUAL, __tag__ MAP<VARCHAR, VARCHAR> METADATA VIRTUAL, proctime as PROCTIME() ) WITH ( 'connector' = 'sls', 'endpoint' ='cn-hangzhou-intranet.log.aliyuncs.com', 'accessId' = '${secret_values.ak_id}', 'accessKey' = '${secret_values.ak_secret}', 'starttime' = '2023-08-30 00:00:00', 'project' ='sls-test', 'logstore' ='sls-input' ); CREATE TEMPORARY TABLE sls_sink( `time` BIGINT, url STRING, dt STRING, float_field FLOAT, double_field DOUBLE, boolean_field BOOLEAN, `__topic__` STRING, `__source__` STRING, `__timestamp__` BIGINT , receive_time BIGINT ) WITH ( 'connector' = 'sls', 'endpoint' ='cn-hangzhou-intranet.log.aliyuncs.com', 'accessId' = '${ak_id}', 'accessKey' = '${ak_secret}', 'project' ='sls-test', 'logstore' ='sls-output' ); INSERT INTO sls_sink SELECT `time`, url, dt, float_field, double_field, boolean_field, `__topic__` , `__source__` , `__timestamp__` , cast(__tag__['__receive_time__'] as bigint) as receive_time FROM sls_input; -
Ingestão de dados com uma fonte de dados SLS
Use o SLS como fonte de dados para ingerir dados em tempo real em sistemas downstream suportados. Por exemplo, a configuração a seguir define um job de ingestão de dados que grava dados de um logstore em um data lake formatado em Paimon no Data Lake Formation (DLF). O job infere automaticamente o schema da tabela de destino e suporta evolução de schema em tempo de execução.
source:
type: sls
name: SLS Source
endpoint: ${endpoint}
project: test_project
logstore: test_log
accessId: ${accessId}
accessKey: ${accessKey}
# Add a primary key to the table.
transform:
- source-table: \.*.\.*
projection: \*
primary-keys: id
# Route all data from test_project.test_log to the test_database.inventory table.
route:
- source-table: test_project.test_log
sink-table: test_database.inventory
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: dlf_uri
catalog.properties.warehouse: your_warehouse
catalog.properties.token.provider: dlf
# (Optional) Enable deletion vectors to improve read performance.
table.properties.deletion-vectors.enabled: true
DataStream API
Para ler ou gravar dados com a DataStream API, use um conector DataStream. Para mais informações, consulte Uso de conectores DataStream.
Se você usar uma versão do VVR anterior à 8.0.10, seu job pode falhar ao iniciar devido a dependências ausentes. Para resolver esse problema, adicione o uber-JAR correspondente como uma dependência adicional.
Leitura do SLS
O Realtime Compute for Apache Flink fornece a classe SlsSourceFunction, uma implementação de SourceFunction, para leitura de dados do Simple Log Service (SLS). O exemplo a seguir lê dados do SLS.
public class SlsDataStreamSource {
public static void main(String[] args) throws Exception {
// Sets up the streaming execution environment
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Creates an SLS source and prints the data to the console.
env.addSource(createSlsSource())
.map(SlsDataStreamSource::convertMessages)
.print();
env.execute("SLS Stream Source");
}
private static SlsSourceFunction createSlsSource() {
SLSAccessInfo accessInfo = new SLSAccessInfo();
accessInfo.setEndpoint("yourEndpoint");
accessInfo.setProjectName("yourProject");
accessInfo.setLogstore("yourLogStore");
accessInfo.setAccessId("yourAccessId");
accessInfo.setAccessKey("yourAccessKey");
// The batch get size is required.
accessInfo.setBatchGetSize(10);
// Optional parameters
accessInfo.setConsumerGroup("yourConsumerGroup");
accessInfo.setMaxRetries(3);
// Start time for consumption, set to the current time.
int startInSec = (int) (new Date().getTime() / 1000);
// Stop time for consumption, where -1 means never stop.
int stopInSec = -1;
return new SlsSourceFunction(accessInfo, startInSec, stopInSec);
}
private static List<String> convertMessages(SourceRecord input) {
List<String> res = new ArrayList<>();
for (FastLogGroup logGroup : input.getLogGroups()) {
int logsCount = logGroup.getLogsCount();
for (int i = 0; i < logsCount; i++) {
FastLog log = logGroup.getLogs(i);
int fieldCount = log.getContentsCount();
for (int idx = 0; idx < fieldCount; idx++) {
FastLogContent f = log.getContents(idx);
res.add(String.format("key: %s, value: %s", f.getKey(), f.getValue()));
}
}
}
return res;
}
}
Gravação no SLS
O Realtime Compute for Apache Flink fornece a classe SLSOutputFormat, uma implementação de OutputFormat, para gravação de dados no SLS. O exemplo a seguir grava dados no SLS.
public class SlsDataStreamSink {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.fromSequence(0, 100)
.map((MapFunction<Long, SinkRecord>) aLong -> getSinkRecord(aLong))
.addSink(createSlsSink())
.name(SlsDataStreamSink.class.getSimpleName());
env.execute("SLS Stream Sink");
}
private static OutputFormatSinkFunction createSlsSink() {
Configuration conf = new Configuration();
conf.setString(SLSOptions.ENDPOINT, "yourEndpoint");
conf.setString(SLSOptions.PROJECT, "yourProject");
conf.setString(SLSOptions.LOGSTORE, "yourLogStore");
conf.setString(SLSOptions.ACCESS_ID, "yourAccessId");
conf.setString(SLSOptions.ACCESS_KEY, "yourAccessKey");
SLSOutputFormat outputFormat = new SLSOutputFormat(conf);
return new OutputFormatSinkFunction<>(outputFormat);
}
private static SinkRecord getSinkRecord(Long seed) {
SinkRecord record = new SinkRecord();
LogItem logItem = new LogItem((int) (System.currentTimeMillis() / 1000));
logItem.PushBack("level", "info");
logItem.PushBack("name", String.valueOf(seed));
logItem.PushBack("message", "it's a test message for " + seed.toString());
record.setContent(logItem);
return record;
}
}
XML
O conector DataStream SLS está disponível no repositório central Maven.
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-sls</artifactId>
<version>${vvr-version}</version>
<exclusions>
<exclusion>
<groupId>org.apache.flink</groupId>
<artifactId>flink-format-common</artifactId>
</exclusion>
</exclusions>
</dependency>