O Realtime Compute for Apache Flink lê dados de log do Simple Log Service ao criar uma tabela source. Este tópico aborda a criação da tabela source e a extração de campos de atributo.
Capacidades suportadas
O Realtime Compute for Apache Flink oferece as seguintes capacidades de consumo de logs.
|
Categoria |
Descrição |
|
Tipo suportado |
Tabela source e tabela result. |
|
Modo de execução |
Apenas modo streaming. |
|
Métrica |
Não suportada. |
|
Formato de dados |
Nenhum. |
|
Tipo de API |
SQL. |
|
Operações na tabela result |
Apenas inserção. Não há suporte para atualizações e exclusões. |
Para começar a consumir dados de log, siga as instruções em Introdução a um deployment Flink SQL.
Pré-requisitos
Se você utilizar um usuário RAM ou uma função RAM, verifique se ele possui as permissões necessárias no console do Realtime Compute for Apache Flink. Para mais informações, consulte Permissões.
Crie um workspace do Realtime Compute for Apache Flink. Para mais informações, consulte Ativar o Realtime Compute for Apache Flink.
Crie um projeto e um Logstore. Consulte Criar um projeto e um Logstore.
Limitações
Apenas o Ververica Runtime (VVR) 11.1 e versões posteriores suportam o uso do SLS como source síncrono para ingestão de dados definida em YAML.
O conector SLS garante apenas a semântica at-least-once.
Não defina o
source parallelismcom um valor superior ao 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.
Criar uma tabela source e uma tabela result
Desenvolva um rascunho SQL completo que contenha uma tabela source e uma tabela result antes de consumir dados de log. Após o processamento, os resultados são inseridos na tabela result por meio da instrução INSERT INTO.
Para mais informações sobre como desenvolver um rascunho SQL no Realtime Compute for Apache Flink, consulte Visão geral do desenvolvimento de jobs.
O Simple Log Service armazena dados de log em tempo real, e o Realtime Compute for Apache Flink os lê em modo streaming. Exemplo de entrada de log:
__source__: 11.85.*.199
__tag__:__receive_time__: 1562125591
__topic__: test-topic
request_method: GET
status: 200
Código de exemplo
O exemplo de rascunho SQL a seguir consome dados do Simple Log Service no Realtime Compute for Apache Flink.
Coloque nomes de tabelas, colunas ou campos reservados entre crases (`) caso haja conflitos entre eles.
CREATE TEMPORARY TABLE sls_input(
request_method STRING,
status BIGINT,
`__topic__` STRING METADATA VIRTUAL,
`__source__` STRING METADATA VIRTUAL,
`__timestamp__` BIGINT 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(
request_method STRING,
status BIGINT,
`__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
request_method,
status,
`__topic__` ,
`__source__` ,
`__timestamp__` ,
cast(__tag__['__receive_time__'] as bigint) as receive_time
FROM sls_input;
Opções WITH
-
Gerais
Parâmetro
Descrição
Tipo
Obrigatório
Padrão
Observações
connector
Conector a ser utilizado.
String
Sim
Nenhum
Defina este valor como
sls.endPoint
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), utilize 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
Nome do projeto SLS.
String
Sim
Nenhum
N/A
logStore
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
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
AccessKey secret da sua conta Alibaba Cloud.
String
Sim
Nenhum
-
-
Específicas da source
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 source adapta-se automaticamente a alterações nos shards e distribui os shards da forma mais uniforme possível entre todas as subtasks da source.
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 tiver shards somente leitura, algumas subtasks 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 a quantidade de shards e simplificar a alocação de tarefas.
shardDiscoveryIntervalMs
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
Modo de inicialização da tabela source.
String
Não
timestamp
-
timestamp(padrão): Consome logs a partir do horário inicial especificado. -
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
Horário inicial para o consumo de logs.
String
Não
Horário 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__do SLS, e não no atributo__timestamp__.stopTime
Horário final 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 um horário no passado. Se você a definir com um horário futuro, 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 será necessário definir
exitAfterFinishcomotrue.
consumerGroup
Nome do grupo de consumidores.
String
Não
Nenhum
Um grupo de consumidores serve para registrar o progresso do 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 utiliza 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 que compartilhem o mesmo grupo de consumidores.
consumeFromCheckpoint
Define se o consumo deve ocorrer 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 de consumidores não possuir um checkpoint correspondente, o consumo começará 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 uma falha na leitura do SLS.
String
Não
3
N/A
batchGetSize
Quantidade de grupos de logs a serem 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
Utilize esta opção para filtrar dados do SLS antes que o Flink os consuma. Isso reduz custos e melhora a velocidade de processamento.
Por exemplo,
'query' = '*| where request_method = ''GET'''indica que, antes de ler dados do SLS, o Flink primeiro corresponde aos dados onde o valor do camporequest_methodé 'GET'.NotaEsta opção utiliza 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 taxas 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 que o Flink os consuma, 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 faça a leitura do SLS.NotaEsta opção utiliza a linguagem SPL do Log Service (SLS). Para mais informações, consulte Sintaxe SPL. Para obter informações sobre como criar ou atualizar um processador SLS, consulte Gerenciar processadores.
ImportanteEsta opção é suportada apenas no VVR 11.3 e posteriores.
Este recurso gera taxas do Log Service (SLS). Para mais informações, consulte Faturamento do Log Service.
-
-
Específicas da sink
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 o horário de gravação do log.String
Não
Horário atual
O valor desta opção deve ser um campo
INTexistente na tabela. Se esta opção não for especificada, o horário atual será usado.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 vão para o 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 de 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
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.
-
Extrair campos de atributo
Além dos campos de log e campos personalizados, o Realtime Compute for Apache Flink pode extrair os seguintes campos de atributo.
|
Campo |
Tipo |
Descrição |
|
__source__ |
STRING METADATA VIRTUAL |
A origem da mensagem. |
|
__topic__ |
STRING METADATA VIRTUAL |
O tópico da mensagem. |
|
__timestamp__ |
BIGINT METADATA VIRTUAL |
O horário do log. |
|
__tag__ |
MAP<VARCHAR, VARCHAR> METADATA VIRTUAL |
A tag da mensagem. Para o atributo |
Para extrair campos de atributo, defina cabeçalhos na sua instrução SQL. Exemplo:
create table sls_stream(
__timestamp__ bigint HEADER,
__receive_time__ bigint HEADER
b int,
c varchar
) with (
'connector' = 'sls',
'endpoint' ='cn-hangzhou.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'
);
Referências
Para mais informações sobre como usar a API DataStream do Realtime Compute for Apache Flink para consumir dados de log, consulte API DataStream.