Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Log Service (SLS)

Última atualização: Jun 30, 2026

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 "__tag__:__receive_time__":"1616742274", '__receive_time__' e '1616742274' são armazenados como um par chave-valor no mapa. Para acessar o valor em SQL, use __tag__['__receive_time__'].

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 parallelism maior 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 de failover, 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.

    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

    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 é 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 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 enableNewSource está definido como true.

    • 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 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

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

    Nota

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

    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.

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

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

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

    Nota

    Esta 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 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 do consumo pelo Flink, 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 leia do SLS.

    Nota

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

    Importante

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

    Nota

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

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
  • Por padrão, o Realtime Compute for Apache Flink não acessa a internet. Use um NAT Gateway para habilitar a comunicação entre sua Virtual Private Cloud (vpc) e a internet. Para mais informações, consulte Como acesso a internet?.

  • Não recomendamos acessar o Log Service (SLS) pela internet. Caso seja necessário, use HTTPS e ative a aceleração de transferência para o SLS.

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

  • continuous: Realiza inferência de schema para cada registro de dados. Se os schemas forem incompatíveis, um schema mais amplo é inferido e um evento de alteração de schema é gerado.

  • static: Realiza inferência de schema apenas uma vez, na inicialização do job. Os dados subsequentes são analisados com base no schema inicial e nenhum evento de alteração de schema é gerado.

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

  • timestamp (padrão): Consome logs a partir de um timestamp específico.

  • latest: Consome logs a partir do offset mais recente.

  • earliest: Consome logs a partir do offset mais antigo.

  • consumer_group: Consome logs a partir do offset registrado no grupo de consumidores. Se nenhum offset estiver registrado para um shard, o consumo começa a partir do offset mais antigo.

startTime

A hora de início para o consumo de logs.

String

Não

Hora atual

O formato é yyyy-MM-dd HH:mm:ss.

Este parâmetro só tem efeito quando startupMode está definido como timestamp.

Nota

Os parâmetros startTime e stopTime baseiam-se no atributo __receive_time__ no Log Service (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

Se desejar que o job Flink encerre após o consumo de todos os logs, também defina exitAfterFinish=true.

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

  • true: O job Flink encerra após o consumo de todos os dados.

  • false (padrão): O job Flink não encerra após o consumo de todos os dados.

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, 'query' = '*| where request_method = ''GET''' filtra dados onde o campo request_method é 'GET' antes que o Flink leia os dados.

Nota

A consulta deve usar a sintaxe SPL do Log Service. Para mais informações, consulte Sintaxe SPL.

Importante
  • Para informações sobre as regiões onde este recurso está disponível no Log Service (SLS), consulte Consumir logs com base em regras.

  • Este recurso está em Beta e é gratuito. Poderão ser aplicadas cobranças futuramente. Para mais informações, consulte Preços.

compressType

O tipo de compressão para o Log Service (SLS).

String

Não

Nenhum

Os tipos de compressão suportados incluem:

  • lz4

  • deflate

  • zstd

timeZone

O fuso horário para startTime e stopTime.

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 __source__, __topic__, __timestamp__ e __tag__. Separe múltiplos campos com vírgula.

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 ,. Por exemplo, se o registro de log SLS upstream for {"col0":"a", "col1":"b", "col2":"c"}, os resultados para diferentes configurações de parâmetros são os seguintes:

Configuração

Table ID

Nenhum

Todas as mensagens são Project.Logstore

col0

a

col0,col1

a.b

col0,col1,col2

a.b.c

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 , para separar múltiplas definições de campos. Por exemplo, id BIGINT, name VARCHAR(10) especifica o tipo do campo id como BIGINT e o tipo do campo name como VARCHAR(10).

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:

  • SQL: Analisa timestamps de entrada no formato yyyy-MM-dd HH:mm:ss.s{precision} (por exemplo, 2020-12-30 12:13:14.123) e os produz no mesmo formato.

  • ISO-8601: Analisa timestamps de entrada no formato yyyy-MM-ddTHH:mm:ss.s{precision} (por exemplo, 2020-12-30T12:13:14.123) e os produz no mesmo formato.

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 ingestion.ignore-errors estiver ativado, o job falhará quando o número acumulado de erros exceder este valor.

Integer

Não

-1

Este parâmetro só tem efeito quando ingestion.ignore-errors está ativado. O valor padrão -1 significa que o job ignora todas as exceções de análise.

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é maxPreFetchLogGroups grupos 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.

    Nota

    Para 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 como continuous, 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

Importante

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>

Perguntas frequentes

Como resolver um OOM do TaskManager (java.lang.OutOfMemoryError: Java heap space) ao restaurar um programa Flink com falha?