Todos os produtos
Search
Central de documentação

Simple Log Service:Consumir dados de log com o Realtime Compute for Apache Flink

Última atualização: Jul 03, 2026

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

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 parallelism com 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 de failover, impedindo o consumo de alguns shards.

Criar uma tabela source e uma tabela result

Importante

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.

Importante

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.

    Importante

    Para 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 é true para 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 consumerGroup para registrar o progresso de consumo no grupo de consumidores do SLS. Em seguida, defina a opção consumeFromCheckpoint como true e 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 enableNewSource está definido como true.

    • 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 consumeFromCheckpoint como true. 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 startupMode está definido como timestamp.

    Nota

    As opções startTime e stopTime baseiam-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 exitAfterFinish como true.

    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.

    Nota

    Jobs 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.

    Importante

    Este parâmetro não é mais suportado no VVR 11.1 e posteriores. Para essas versões, defina a opção startupMode como consumer_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 batchGetSize nã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

    Importante

    Esta 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 campo request_method é 'GET'.

    Nota

    Esta 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 query quanto processor forem especificados, query terá precedência e processor será 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 processor em vez de query.

    Por exemplo, 'processor' = 'test-filter-processor' indica que o processador SLS filtra os dados antes que o Flink faça a leitura do SLS.

    Nota

    Esta 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.

    Importante

    Esta 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 INT existente 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 partitionField for 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.

    Nota

    Esta 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 "__tag__:__receive_time__":"1616742274", os campos __receive_time__ e 1616742274 são registrados como pares chave-valor em um mapa. Você pode incluir __tag__['__receive_time__'] em uma instrução SQL para consultar a tag.

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.