O Flink SQL integra-se à Simple Log Service Processing Language (SPL) para analisar dados de log semiestruturados em campos estruturados destinados à análise, sem criar logstores intermediários ou tabelas temporárias.
Contexto
O Simple Log Service (SLS) é uma plataforma nativa da cloud de observabilidade e análise que oferece services em larga escala, de baixo custo e em tempo real para logs, métricas e rastreamentos. O Realtime Compute for Apache Flink é uma plataforma de análise de big data desenvolvida pela Alibaba Cloud com base no Apache Flink, amplamente utilizada para análise de dados em tempo real e monitoramento de riscos. O Realtime Compute for Apache Flink suporta nativamente o conector do SLS, permitindo usar o SLS como tabela de source ou de resultados.
O conector do SLS para o Realtime Compute for Apache Flink processa logs estruturados diretamente ao mapear os campos de log um a um com as colunas da tabela do Flink SQL. No entanto, muitos logs de negócios não são totalmente estruturados. Por exemplo, todo o conteúdo do log pode ser gravado em um único campo, exigindo expressões regulares ou divisão por delimitadores para extrair campos estruturados. Com a SPL no conector do SLS, é possível estruturar esses dados por meio da limpeza de logs e da normalização de formatos.
Dados de log semiestruturados
O log de exemplo a seguir possui um formato complexo que mistura strings JSON com outros tipos de dados. O log contém:
O campo Payload, uma string JSON em que o campo aninhado schedule também é uma estrutura JSON.
O campo requestURL, que representa um caminho de URL padrão.
O campo error, que começa com a string
CouldNotExecuteQueryseguida por uma estrutura JSON.O campo __tag__:__path__, que contém o caminho do arquivo de log, onde
service_apode representar o nome do service.O campo caller, que inclui o nome do arquivo e o número da linha.
{
"Payload": "{\"lastNotified\": 1705030483, \"serverUri\": \"http://test.alert.com/alert-api/tasks\", \"jobID\": \"44d6ce47bb4995ef0c8052a9a30ed6d8\", \"alertName\": \"alert-12345678-123456\", \"project\": \"test-sls-project\", \"projectId\": 123, \"aliuid\": \"1234567890\", \"alertDisplayName\": \"\\u6d4b\\u8bd5\\u963f\\u91cc\\u4e91\\u544a\\u8b66\", \"checkJobUri\": \"http://test.alert.com/alert-api/task_check\", \"schedule\": {\"timeZone\": \"\", \"delay\": 0, \"runImmediately\": false, \"type\": \"FixedRate\", \"interval\": \"1m\"}, \"jobRunID\": \"bf86aa5e67a6891d-61016da98c79b-5071a6b\", \"firedNotNotified\": 25161}",
"TaskID": "bf86aa5e67a6891d-61016da98c79b-5071a6b-334f81a-5c38aaa1-9354-43ec-8369-4f41a7c23887",
"TaskType": "ALERT",
"__source__": "11.199.XXX.XXX",
"__tag__:__hostname__": "iabcde12345.cloud.abc121",
"__tag__:__path__": "/var/log/service_a.LOG",
"caller": "executor/pool.go:64",
"error": "CouldNotExecuteQuery : {\n \"httpCode\": 404,\n \"errorCode\": \"LogStoreNotExist\",\n \"errorMessage\": \"logstore k8s-event does not exist\",\n \"requestID\": \"65B7C10AB43D9895A8C3DB6A\"\n}",
"requestURL": "/apis/autoscaling/v2beta1/namespaces/python-etl/horizontalpodautoscalers/cn-shenzhen-56492-1234567890123?timeout=30s",
"ts": "2024-01-29 22:57:13"
}
Requisitos de estruturação de dados
Para extrair informações valiosas desses logs, transforme os dados extraindo os campos principais para análise.
Do campo error, extraia httpCode, errorCode, errorMessage e requestID.
Do campo __tag__:__path_, extraia _service_a como serviceName.
Do campo caller, extraia pool.go como fileName e 64 como fileNo.
Do campo Payload, extraia project. Do objeto aninhado schedule dentro de Payload, extraia type e nomeie-o como scheduleType.
Renomeie o campo __source__ para serviceIP.
Todos os demais campos são descartados. A lista final de campos necessários é: httpCode, errorCode, errorMessage, requestID, serviceName, fileName, fileNo, project, scheduleType e serviceIP.
Soluções
Existem vários métodos para transformar os dados. As soluções a seguir utilizam SLS e Flink, cada uma adequada a diferentes cenários.
Solução de transformação de dados: No console do SLS, crie um job de transformação de dados para limpar os dados e armazená-los em um logstore de destino.
Solução Flink: Defina
errorepayloadcomo campos na tabela de source. Utilize as funções de expressão regular e JSON do Flink SQL para analisar esses campos, grave os resultados em uma tabela temporária e, em seguida, realize a análise nessa tabela.Solução SPL: Configure uma instrução SPL no conector SLS do Flink para transformar os dados. A tabela de source do Flink é então definida com o esquema final estruturado.
A abordagem SPL é mais leve: evita a criação de um logstore intermediário (necessário na solução de transformação de dados) e de uma tabela temporária no Flink (necessária na solução Flink). Ao realizar a transformação de dados mais próximo da source, você foca na lógica de negócios na plataforma de computação, criando uma separação de responsabilidades mais clara.
Uso da SPL no Flink
1. Preparar dados no SLS
Certifique-se de ter ativado o Simple Log Service e criado um projeto e um logstore.
-
Grave o trecho de log anterior no logstore de destino usando um SDK para simular dados de amostra.
Após a gravação dos dados, acesse a aba Raw Logs no console do Simple Log Service para visualizar os dados simulados. O conteúdo do log inclui campos como TaskType (com o valor
ALERT), Payload (que contémalertName,projectealiuid) e informações de erro (código de erroLogStoreNotExist, mensagemlogstore k8s-event does not existe status HTTP 404). -
No logstore, escreva uma instrução de pipeline SPL do SLS e visualize os resultados.
Após executar a instrução SPL, a aba Raw Logs exibe os registros de log analisados, que incluem campos estruturados como
errorCode,errorMessage,fileName,fileNo,httpCode,project,requestID,scheduleTypeeserviceHost.A instrução de consulta SPL é apresentada a seguir. A sintaxe de pipeline SPL usa o delimitador pipe (|) para separar comandos. Insira um comando por vez para ver o resultado imediato e adicione mais pipes para construir iterativamente a consulta final. Para obter mais informações, consulte Sintaxe de consulta de varredura.
* | project Payload, error, "__tag__:__path__", "__tag__:__hostname__", caller | parse-json Payload | project-away Payload | parse-regexp error, 'CouldNotExecuteQuery : ({[\w":\s,\-}]+)' as errorJson | parse-json errorJson | parse-regexp "__tag__:__path__", '\/var\/log\/([\w\_]+).LOG' as serviceName | parse-regexp caller, '\w+/([\w\.]+):(\d+)' as fileName, fileNo | project-rename serviceHost="__tag__:__hostname__" | extend scheduleType = json_extract_scalar(schedule, '$.type') | project httpCode, errorCode,errorMessage,requestID,fileName, fileNo, serviceHost,scheduleType, projectExplicação da sintaxe:
Linha 1: O comando project retém os campos Payload, error, __tag__:__path__ e caller para análise e descarta todos os outros campos.
Linha 2: O comando parse-json expande a string Payload em um objeto JSON. Seus campos de nível superior, como lastNotified, serviceUri e jobID, são adicionados ao resultado.
Linha 3: O comando project-away remove o campo Payload original.
Linha 4: O comando parse-regexp utiliza uma expressão regular para extrair a parte JSON do campo error e a atribui a um novo campo chamado errorJson.
Linha 5: O comando parse-json expande o campo errorJson, extraindo httpCode, errorCode e errorMessage.
Linha 6: O comando parse-regexp usa uma expressão regular para extrair o nome do arquivo de __tag__:__path__ e o nomeia como serviceName.
Linha 7: O comando parse-regexp emprega uma expressão regular para extrair o nome do arquivo e o número da linha do campo caller, atribuindo-os aos campos fileName e fileNo.
Linha 8: O comando project-rename renomeia o campo __tag__:__hostname__ para serviceHost.
Linha 9: O comando extend utiliza a função
json_extract_scalarpara extrair o campotypedo objeto schedule e o nomeia como scheduleType.Linha 10: O comando project retém apenas os campos finais necessários, incluindo o campo project extraído de Payload.
2. Criar um job SQL
Faça login no console do Realtime Compute for Apache Flink e clique em no workspace de destino.
No painel de navegação à esquerda, escolha .
Clique em Create. Na caixa de diálogo New Draft, escolha e clique em Next.
-
No editor de rascunho, insira a seguinte instrução para criar uma tabela temporária.
CREATE TEMPORARY TABLE sls_input_complex ( errorCode STRING, errorMessage STRING, fileName STRING, fileNo STRING, httpCode STRING, requestID STRING, scheduleType STRING, serviceHost STRING, project STRING, proctime as PROCTIME() ) WITH ( 'connector' = 'sls', 'endpoint' ='cn-beijing-intranet.log.aliyuncs.com', 'accessId' = '${yourAccessKeyID}', 'accessKey' = '${yourAccessKeySecret}', 'starttime' = '2024-02-01 10:30:00', 'project' ='${project}', 'logstore' ='${logtore}', 'query' = '* | project Payload, error, "__tag__:__path__", "__tag__:__hostname__", caller | parse-json Payload | project-away Payload | parse-regexp error, ''CouldNotExecuteQuery : ({[\w":\s,\-}]+)'' as errorJson | parse-json errorJson | parse-regexp "__tag__:__path__", ''\/var\/log\/([\w\_]+).LOG'' as serviceName | parse-regexp caller, ''\w+/([\w\.]+):(\d+)'' as fileName, fileNo | project-rename serviceHost="__tag__:__hostname__" | extend scheduleType = json_extract_scalar(schedule, ''$.type'') | project httpCode, errorCode,errorMessage,requestID,fileName, fileNo, serviceHost,scheduleType,project' );A tabela a seguir descreve os parâmetros na cláusula WITH. Substitua os valores de exemplo pelos seus valores reais.
Parâmetro
Descrição
Exemplo
connector
O tipo de conector. Conectores suportados.
sls
endpoint
O 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
Hora de início para consumo de logs.
2025-02-19 00:00:00
project
O nome do projeto SLS.
test-project
logstore
O nome do Logstore do SLS.
clb-access-log
query
A instrução SPL. Escape literais de string com aspas simples duplicadas (
'') no Flink SQL.*
where slbid = ''slb-01''
-
Selecione a instrução SQL, clique em com o botão direito e selecione Running para conectar-se ao Simple Log Service.
CREATE TEMPORARY TABLE sls_input_complex ( errorCode STRING, errorMessage STRING, fileName STRING, fileNo STRING, httpCode STRING, requestID STRING, scheduleType STRING, serviceHost STRING, project STRING, proctime as PROCTIME() ) WITH ( 'connector' = 'sls', 'endpoint' = 'cn-hxxx', 'accessId' = 'xxx', 'accessKey' = 'xxx', 'starttime' = 'xxx', 'project' = 'xxx', 'logstore' = 'clb7xxx', 'query' = '* | prxxx", "__tag__:__hosxxx' );
3. Executar consulta e visualizar resultados
-
No editor de jobs, insira a seguinte instrução para consultar os dados:
SELECT * FROM sls_input_complex; -
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 abaixo.
Defina Name como
demo-test, Deployment Target comodefault-queue, Status como RUNNING e Engine Version comovvr-8.0.11-flink-1.17, e clique em Create Session Cluster. -
Na caixa de diálogo de depuração, selecione o cluster de sessão criado e clique em OK.
Nota: A depuração com uma tabela de source SLS avança o deslocamento do grupo de consumidores. Um job implantado será retomado a partir desse novo deslocamento.
-
Na aba Results, observe que cada coluna na tabela corresponde a um campo processado pela consulta SPL.
Após executar o Debug da consulta, a tabela de resultados mostra os campos analisados. Por exemplo, a coluna
errorCodeapresentaLogStoreNotExist, a colunaerrorMessageexibelogstore k8s-event does not existe a colunahttpCodemostra404.