Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Log Service (SLS)

Última atualização: Aug 13, 2026

Saiba como usar o conector do Log Service (SLS).

Contexto

O Simple Log Service é um serviço de ponta a ponta para dados de log. Ele auxilia na coleta, no consumo, no envio, na consulta e na análise eficiente desses dados, melhorando a eficiência de O&M e permitindo processar grandes volumes de informações.

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 são apenas de adição (append-only); não é possível atualizar ou excluir dados.

Recursos

O conector de source SLS lê diretamente os campos de atributos das mensagens. A tabela abaixo detalha 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

Horário 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 projeto do Log Service e um Logstore. Para mais informações, consulte Create a Project and a 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 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 na quantidade 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

    Conector a ser usado.

    String

    Sim

    Nenhum

    Defina 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 Service endpoints.

    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 detalhes, veja How do I access the Internet?.

    • Recomendamos não acessar o SLS pela internet. Caso seja necessário, use HTTPS e ative a aceleração de transferência. Consulte Manage transfer acceleration para saber mais.

    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 Obtain an AccessKey pair.

    Importante

    Para evitar a exposição do seu par AccessKey, recomendamos usar variáveis para especificar o AccessKey ID e o AccessKey secret. Veja Project variables para detalhes.

    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 fonte adapta-se automaticamente a mudanças nos shards e distribui-os da forma mais equilibrada possível entre todas as subtarefas de source.

    Importante
    • Esta opção é suportada apenas no VVR 8.0.9 e posteriores. O valor padrão é true para VVR 11,1 e superiores.

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

    • Caso um Logstore possua shards somente leitura, algumas subtarefas podem continuar solicitando dados de outros shards após concluírem os seus próprios. Isso pode gerar desequilíbrio na carga de trabalho e afetar o desempenho. Para mitigar esse problema, ajuste o paralelismo, otimize a estratégia de agendamento ou mescle shards pequenos para reduzir a quantidade total 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.

    • Suportado apenas no VVR 8.0.9 e versões posteriores.

    startupMode

    Modo de inicialização da tabela de source.

    String

    Não

    timestamp

    • timestamp (padrão): Consome logs a partir do horário inicial especificado.

    • latest: Inicia o consumo de logs pelo offset mais recente.

    • earliest: Começa o consumo de logs pelo offset mais antigo.

    • consumer_group: Inicia o consumo a partir do offset registrado pelo grupo de consumidores. Se o grupo não tiver um offset registrado para determinado shard, o consumo começa pelo 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 no 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 consumo de logs.

    String

    Não

    Horário atual

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

    Tem efeito apenas 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 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 configurada com uma data passada. Se definida para um horário futuro, o consumo poderá parar prematuramente caso nenhum novo log seja gravado, causando interrupção no fluxo de dados sem mensagens de erro.

    • Para que o job Flink encerre após o consumo de todos os logs, também é necessário definir exitAfterFinish como true.

    consumerGroup

    Nome do grupo de consumidores.

    String

    Não

    Nenhum

    Um grupo de consumidores registra o progresso do consumo. Você pode definir um nome personalizado sem formato fixo.

    Nota

    Jobs Flink diferentes devem usar grupos de consumidores distintos. Se múltiplos jobs utilizarem o mesmo grupo, eles não se coordenarão e cada um 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. Consequentemente, cada consumidor processa mensagens independentemente, mesmo compartilhando o mesmo grupo.

    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 dos logs a partir do checkpoint salvo nesse grupo. Caso não exista um checkpoint correspondente, o consumo começa pelo 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 versões 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 por requisiçã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 todo o dado ser consumido.

    String

    Não

    false

    • true: O programa Flink encerra após o consumo completo dos 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 selecionará apenas registros onde o campo request_method tenha o valor 'GET'.

    Nota

    Esta opção utiliza a linguagem SPL do Log Service (SLS). Para mais informações, consulte SPL syntax.

    Importante
    • Suportado apenas no VVR 8.0.1 e versões posteriores.

    • Este recurso gera cobranças do Log Service (SLS). Consulte Billing of Log Service para detalhes.

    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, reduzindo custos e melhorando a velocidade de processamento. Recomendamos o uso de processor em vez de query.

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

    Nota

    Esta opção utiliza a linguagem SPL do Log Service (SLS). Para mais informações, consulte SPL syntax. Para saber como criar ou atualizar um processador SLS, veja Manage processors.

    Importante

    Suportado apenas no VVR 11.3 e versões posteriores.

    Este recurso gera cobranças do Log Service (SLS). Consulte Billing of Log Service para detalhes.

    preserveRawBytes

    Define se os bytes brutos transportados pelo SLS devem ser buscados e preservados diretamente quando enableNewSource estiver definido como true.

    Boolean

    Não

    false

    Quando esta opção está ativada, campos BINARY/VARBINARY leem o byte[] bruto diretamente, e campos CHAR/VARCHAR constroem strings a partir dos bytes brutos. Outros campos continuam usando a lógica de conversão padrão. O comportamento é consistente com o da fonte antiga.

    Nota
    • Suportado apenas no VVR 11.8 e versões posteriores.

    • Esta opção só tem efeito quando enableNewSource está definido como true.

  • 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__, indicando 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__, indicando 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 não especificado, o horário atual será usado.

    sourceField

    Especifica um campo cujo valor sobrescreve o atributo __source__, indicando a source do log, como o endereço IP da máquina que o gerou.

    String

    Não

    Nenhum

    O valor desta opção deve ser um campo existente na tabela.

    partitionField

    Define um campo para particionamento. Um hash do valor deste campo determina qual shard receberá 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 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

    Suportado apenas no VVR 8.0.6 e versões posteriores.

Mapeamento de tipos

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 pelo Realtime Compute for Apache Flink versão 11,1 e posteriores.

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

Tipo da fonte de dados.

String

Sim

Nenhum

O valor deve ser sls.

endpoint

Endpoint do Log Service (SLS).

String

Sim

Nenhum

Endereço de acesso VPC do Log Service (SLS). Para mais informações, consulte Service Endpoints.

Nota
  • Por padrão, o Realtime Compute for Apache Flink não acessa a internet. Use um NAT Gateway para permitir a comunicação entre sua Virtual Private Cloud (VPC) e a internet. Consulte How can I access the internet? para detalhes.

  • Não recomendamos acessar o Log Service (SLS) pela internet. Se for indispensável, use HTTPS e ative a transfer acceleration para o SLS.

accessId

AccessKey ID da sua conta Alibaba Cloud.

String

Sim

Nenhum

Para mais informações, consulte How do I view the AccessKey ID and AccessKey secret information?.

Importante

Para evitar a exposição das suas informações AccessKey, recomendamos usar uma variável de projeto para especificar o valor do AccessKey. Veja Project variables para mais detalhes.

accessKey

AccessKey secret da sua conta Alibaba Cloud.

String

Sim

Nenhum

project

Nome do projeto Log Service (SLS).

String

Sim

Nenhum

Nenhum

logStore

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

Estratégia de inferência de schema.

String

Não

continuous

  • continuous: Realiza inferência de schema para cada registro de dados. Se houver incompatibilidade, um schema mais amplo é inferido e um evento de alteração de schema é gerado.

  • static: Realiza a inferência de schema apenas uma vez, na inicialização do job. Os dados subsequentes são analisados com base no schema inicial, sem gerar eventos de alteração.

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 uma quantidade especificada 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

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 inicia pelo offset mais antigo.

startTime

Horário inicial para consumo de logs.

String

Não

Horário 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__ do Log Service (SLS), e não no atributo __timestamp__.

stopTime

Horário final para consumo de logs.

String

Não

Nenhum

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

Nota

Para que o job Flink encerre após o consumo de todos os logs, também é necessário definir exitAfterFinish=true.

consumerGroup

Nome do grupo de consumidores.

String

Não

Nenhum

Um grupo de consumidores registra o progresso do consumo. Você pode especificar qualquer nome personalizado.

batchGetSize

Número de grupos de logs lidos por requisiçã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 em caso de falha na leitura do Log Service (SLS).

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 completo dos dados.

  • false (padrão): O job Flink não encerra após o consumo completo dos 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, economizando custos e aumentando a velocidade de processamento.

Por exemplo, 'query' = '*| where request_method = ''GET''' filtra dados onde o campo request_method é 'GET' antes que sejam lidos pelo Flink.

Nota

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

Importante
  • Para informações sobre as regiões onde este recurso está disponível no Log Service (SLS), consulte Consume logs based on rules.

  • Este recurso está em Beta e é gratuito. Poderá haver cobrança futuramente. Para mais detalhes, veja Pricing.

compressType

Tipo de compressão para o Log Service (SLS).

String

Não

Nenhum

Tipos de compressão suportados incluem:

  • lz4

  • deflate

  • zstd

timeZone

Fuso horário para startTime e stopTime.

String

Não

Nenhum

Por padrão, nenhum deslocamento é adicionado.

regionId

Região onde o Log Service (SLS) está implantado.

String

Não

Nenhum

Para mais informações, consulte Supported regions.

signVersion

Versão da assinatura da requisição para o Log Service (SLS).

String

Não

Nenhum

Para mais informações, consulte Request signatures.

shardModDivisor

Divisor usado na leitura de shards do Logstore.

Int

Não

-1

Para mais informações, consulte Shards.

shardModRemainder

Resto usado na leitura de shards do Logstore.

Int

Não

-1

Para mais informações, consulte Shards.

metadata.list

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:

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 no Realtime Compute for Apache Flink versão 11,6 e posteriores.

fixed-types

Especifica os tipos de dados para campos específicos ao analisar logs 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) define o tipo do campo id como BIGINT e o tipo do campo name como VARCHAR(10).

Nota

Este parâmetro é suportado no Realtime Compute for Apache Flink versão 11,6 e posteriores.

timestamp-format.standard

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 gera saída 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 gera saída no mesmo formato.

Nota

Este parâmetro é suportado no Realtime Compute for Apache Flink versão 11,6 e posteriores.

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 no Realtime Compute for Apache Flink versão 11,6 e posteriores.

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 no Realtime Compute for Apache Flink versão 11,6 e posteriores.

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 elimina 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 utilizando 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 os dados, 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 nesse 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 do horário atual para analisar o schema do log.

  • Informações de chave primária

    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, infere as colunas físicas e as compara com o schema atual. Se o schema inferido for inconsistente com o atual, os schemas são mesclados conforme 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 tiverem campos que existem no schema atual, o conector retém esses campos, preenche seus dados com NULL e não gera 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 source 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 source SLS

    Use o SLS como fonte de dados para ingerir informações em tempo real em sistemas downstream suportados. Por exemplo, a configuração abaixo define um job de ingestão 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 Usage of DataStream connectors.

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

How to resolve a TaskManager OOM (java.lang.OutOfMemoryError: Java heap space) when restoring a failed Flink program?