Todos os produtos
Search
Central de documentação

Simple Log Service:Flink SQL: Row filtering and column pruning using SPL

Última atualização: Jul 03, 2026

A SPL permite aplicar filtragem de linhas e projeção de colunas diretamente no servidor do Simple Log Service (SLS) antes que os dados cheguem ao Flink, o que reduz a sobrecarga de rede e os custos de computação.

Contexto

Ao usar o SLS como tabela de origem no Realtime Compute for Apache Flink, o conector consome todos os registros do Logstore após um horário inicial especificado. Esse comportamento gera dois problemas:

  1. O conector busca linhas ou colunas em excesso da origem, aumentando a sobrecarga de rede.

  2. A filtragem de dados desnecessários no Flink desperdiça recursos computacionais.

A SPL resolve essa questão com o pushdown de filtros e de projeções para o conector do SLS. Configure o parâmetro query para enviar as condições de filtro e os campos projetados ao servidor, evitando a transferência e o processamento de todo o conjunto de dados.

Funcionamento

  • Sem instrução SPL: O Flink busca todas as linhas e colunas do SLS para processamento.

    image
  • Com instrução SPL: Se a instrução SPL incluir filtragem de linhas ou projeção de colunas, o Flink recupera apenas o subconjunto correspondente.

    image

Pré-requisitos

  • Você já ativou o Simple Log Service e criou um Project e um Logstore.

  • Este exemplo usa um Logstore preenchido com logs simulados de SLB de camada 7. Os dados contêm mais de 10 campos e são gerados continuamente. Exemplo de entrada de log:

    {
      "__source__": "127.0.0.1",
      "__tag__:__receive_time__": "1706531737",
      "__time__": "1706531727",
      "__topic__": "slb_layer7",
      "body_bytes_sent": "3577",
      "client_ip": "114.137.XXX.XXX",
      "host": "www.pi.mock.com",
      "http_host": "www.cwj.mock.com",
      "http_user_agent": "Mozilla/5.0 (Windows NT 6.2; rv:22.0) Gecko/20130405 Firefox/23.0",
      "request_length": "1662",
      "request_method": "GET",
      "request_time": "31",
      "request_uri": "/request/path-0/file-3",
      "scheme": "https",
      "slbid": "slb-02",
      "status": "200",
      "upstream_addr": "42.63.XXX.XXX",
      "upstream_response_time": "32",
      "upstream_status": "200",
      "vip_addr": "223.18.XX.XXX"
    }
  • O campo slbid tem três valores distintos. Uma consulta de 15 minutos mostra que slb-01 e slb-02 apresentam contagens semelhantes.

    Execute * | select slbid, count(1) cnt group by slbid no editor de consultas. Resultados: slb-02 tem 2.001 entradas, slb-01 tem 1.986 entradas e lb-uf6qfcvsrouodhvw39oog tem 4.144 entradas.

Procedimento

Filtragem de linhas: envie as condições de filtro ao servidor do SLS pelo parâmetro query para impedir a transferência completa do conjunto de dados.

Projeção de colunas: envie os campos projetados ao servidor do SLS pelo parâmetro query para retornar apenas as colunas especificadas.

Filtragem de linhas

Etapa 1: Criar um job SQL

  1. Faça login no console do Realtime Compute for Apache Flink e clique em no workspace desejado.

  2. No painel de navegação à esquerda, escolha Data Development > ETL.

  3. Clique em Create. Na caixa de diálogo New Draft, escolha SQL Scripts > Blank Stream Draft e clique em Next.

  4. Copie o seguinte SQL de tabela temporária para o editor.

    CREATE TEMPORARY TABLE sls_input(
      request_uri STRING,
      scheme STRING,
      slbid STRING,
      status STRING,
      `__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' = 'yourAccessKeyID',
      'accessKey' = 'yourAccessKeySecret',
      'starttime' = '2025-02-19 00:00:00',
      'project' ='test-project',
      'logstore' ='clb-access-log',
      'query' = '* | where slbid = ''slb-01'''
    );

    Parâmetros:

    Parâmetro

    Descrição

    Exemplo

    connector

    Tipo do conector. Conectores suportados.

    sls

    endpoint

    Endpoint interno do SLS. Endpoints.

    cn-hangzhou-intranet.log.aliyuncs.com

    accessId

    Seu AccessKey ID. Criar um AccessKey.

    LTAI

    accessKey

    O AccessKey Secret. Criar um AccessKey.

    yourAccessKeySecret

    starttime

    Horário inicial para consumo de logs.

    2025-02-19 00:00:00

    project

    Nome do projeto do SLS.

    test-project

    logstore

    Nome do Logstore do SLS.

    clb-access-log

    query

    Instrução SPL. No Flink SQL, escape literais de string duplicando as aspas simples ('').

    * \

    where slbid = ''slb-01''

  5. Selecione a instrução SQL, clique com o botão direito e selecione Running para se conectar ao SLS.

    CREATE TEMPORARY TABLE sls_input(
      request_uri STRING,
      scheme STRING,
      slbid STRING,
      status STRING,
      `__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-haxxx',
      'accessId' = 'xxx',
      'accessKey' = 'xxx',
      'starttime' = '202xxx',
      'project' = 'test-pxxx',
      'logstore' = 'clb-axxx',
      'query' = '*|whexxx'
    );

Etapa 2: Executar uma consulta e visualizar os resultados

  1. Execute a seguinte consulta de agregação no campo slbid.

    SELECT slbid, count(1) as slb_cnt FROM sls_input GROUP BY slbid;
  2. Clique em Debug no canto superior direito. Na caixa de diálogo, selecione Create new session cluster na lista suspensa Session Cluster e configure conforme descrito abaixo.

    Defina Name como demo-test, Deployment Target como default-queue, Status como RUNNING e Engine Version como vvr-8.0.11-flink-1.17. Em seguida, clique em Create Session Cluster.

  3. Na caixa de diálogo de depuração, selecione o session cluster recém-criado e clique em OK.

    Nota: A depuração com uma tabela de origem do SLS avança o offset do grupo de consumidores. Um job implantado retoma a partir desse novo offset.

  4. Os resultados mostram que slbid é sempre slb-01, confirmando que sls_input contém apenas linhas correspondentes à condição de filtro.

    O valor de slb_cnt é 185.

Projeção de colunas

Etapa 1: Criar um job SQL

  1. Faça login no console do Realtime Compute for Apache Flink e clique em no workspace desejado.

  2. No painel de navegação à esquerda, escolha Data Development > ETL.

  3. Clique em Create. Na caixa de diálogo New Draft, escolha SQL Scripts > Blank Stream Draft e clique em Next.

  4. Copie o seguinte SQL de tabela temporária para o editor. Diferentemente da filtragem de linhas, o parâmetro query adiciona uma instrução de projeção. Os comandos são encadeados com | (semelhante a pipes Unix), e apenas os campos projetados são recuperados do servidor do SLS.

    CREATE TEMPORARY TABLE sls_input_project(
      request_uri STRING,
      scheme STRING,
      slbid STRING,
      status STRING,
      `__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' = 'yourAccessKeyID',
      'accessKey' = 'yourAccessKeySecret',
      'starttime' = '2025-02-19 00:00:00',
      'project' ='test-project',
      'logstore' ='clb-access-log',
      'query' = '* | where slbid = ''slb-01'' | project request_uri, scheme, slbid, status, __topic__, __source__, "__tag__:__receive_time__"'
    );

    Parâmetros:

    Parâmetro

    Descrição

    Exemplo

    connector

    Tipo do conector. Conectores suportados.

    sls

    endpoint

    Endpoint interno do SLS. Endpoints.

    cn-hangzhou-intranet.log.aliyuncs.com

    accessId

    Seu AccessKey ID. Criar um AccessKey.

    LTAI

    accessKey

    O AccessKey Secret. Criar um AccessKey.

    yourAccessKeySecret

    starttime

    Horário inicial para consumo de logs.

    2025-02-19 00:00:00

    project

    Nome do projeto do SLS.

    test-project

    logstore

    Nome do Logstore do SLS.

    clb-access-log

    query

    Instrução SPL. No Flink SQL, escape literais de string duplicando as aspas simples ('').

    * \

    where slbid = ''slb-01''

  5. Selecione a instrução SQL, clique com o botão direito e selecione Running para se conectar ao SLS.

Etapa 2: Executar uma consulta e visualizar os resultados

  1. Execute a seguinte consulta de agregação no campo slbid.

    SELECT slbid, count(1) as slb_cnt FROM sls_input_project GROUP BY slbid;
  2. Clique em Debug no canto superior direito. Na caixa de diálogo, selecione Create new session cluster na lista suspensa Session Cluster e configure conforme descrito abaixo.

    Defina Name como demo-test, Deployment Target como default-queue, Status como RUNNING e Engine Version como vvr-8.0.11-flink-1.17. Em seguida, clique em Create Session Cluster.

  3. Na caixa de diálogo de depuração, selecione o session cluster recém-criado e clique em OK.

    Nota: A depuração com uma tabela de origem do SLS avança o offset do grupo de consumidores. Um job implantado retoma a partir desse novo offset.

  4. Os resultados são semelhantes aos da filtragem de linhas.

    Nota

    Diferentemente da filtragem de linhas, a projeção de colunas retorna apenas os campos especificados, reduzindo ainda mais o tráfego de rede.

    O valor de slb_cnt é 185.