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:
O conector busca linhas ou colunas em excesso da origem, aumentando a sobrecarga de rede.
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.
-
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.
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
slbidtem três valores distintos. Uma consulta de 15 minutos mostra queslb-01eslb-02apresentam contagens semelhantes.Execute
* | select slbid, count(1) cnt group by slbidno editor de consultas. Resultados:slb-02tem 2.001 entradas,slb-01tem 1.986 entradas elb-uf6qfcvsrouodhvw39oogtem 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
Faça login no console do Realtime Compute for Apache Flink e clique em no workspace desejado.
No painel de navegação à esquerda, escolha .
Clique em Create. Na caixa de diálogo New Draft, escolha e clique em Next.
-
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''
-
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
-
Execute a seguinte consulta de agregação no campo
slbid.SELECT slbid, count(1) as slb_cnt FROM sls_input GROUP BY slbid; -
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 comodefault-queue, Status como RUNNING e Engine Version comovvr-8.0.11-flink-1.17. Em seguida, clique em Create Session Cluster. -
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.
-
Os resultados mostram que slbid é sempre
slb-01, confirmando quesls_inputcontém apenas linhas correspondentes à condição de filtro.O valor de slb_cnt é
185.
Projeção de colunas
Etapa 1: Criar um job SQL
Faça login no console do Realtime Compute for Apache Flink e clique em no workspace desejado.
No painel de navegação à esquerda, escolha .
Clique em Create. Na caixa de diálogo New Draft, escolha e clique em Next.
-
Copie o seguinte SQL de tabela temporária para o editor. Diferentemente da filtragem de linhas, o parâmetro
queryadiciona 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''
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
-
Execute a seguinte consulta de agregação no campo
slbid.SELECT slbid, count(1) as slb_cnt FROM sls_input_project GROUP BY slbid; -
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 comodefault-queue, Status como RUNNING e Engine Version comovvr-8.0.11-flink-1.17. Em seguida, clique em Create Session Cluster. -
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.
-
Os resultados são semelhantes aos da filtragem de linhas.
NotaDiferentemente 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.