Este tutorial demonstra como criar uma solução unificada de análise em lote e em tempo real com dados de eventos do GitHub. Você usará o MaxCompute para construir um data warehouse em lote, além do Realtime Compute for Apache Flink e do Hologres para criar um data warehouse em tempo real. Em seguida, utilizará o Hologres e o MaxCompute para executar análises de dados em tempo real e em lote.
Contexto
Com a crescente digitalização dos negócios, aumenta também a demanda por dados mais recentes. Além do processamento tradicional em lote para grandes volumes, muitas empresas precisam lidar com processamento, armazenamento e análise de dados em tempo real. A análise unificada em lote e em tempo real surgiu para atender a essa necessidade.
Essa abordagem gerencia e processa dados em tempo real e em lote em uma única plataforma. Ela permite uma conexão perfeita entre o processamento em tempo real e a análise em lote, melhorando a eficiência e a precisão da análise de dados. Os principais benefícios incluem:
Maior eficiência no processamento de dados: A integração de dados em tempo real e em lote em uma única plataforma aprimora significativamente a eficiência do processamento e reduz os custos de transferência e conversão de dados.
Melhor precisão nas análises: A combinação de dados em tempo real e em lote para análise aumenta a precisão e a exatidão dos resultados.
Simplificação da gestão e do processamento de dados, tornando essas tarefas mais eficientes.
Aproveitamento total do valor dos dados para oferecer melhor suporte à tomada de decisões empresariais.
O Alibaba Cloud oferece uma solução simplificada e unificada de data warehouse para cenários em lote e em tempo real. Essa solução utiliza o MaxCompute para processamento em lote e o Hologres para análise em tempo real. Combinados com as capacidades de processamento em tempo real do Realtime Compute for Apache Flink, esses serviços formam o motor principal do data warehouse unificado do Alibaba Cloud.
Arquitetura da solução
O diagrama a seguir ilustra o pipeline completo para implementar a análise unificada em lote e em tempo real no conjunto de dados público de eventos do GitHub, utilizando MaxCompute e Hologres.

Nesta arquitetura, uma instância ECS coleta e agrega dados de eventos em tempo real e em lote do GitHub, servindo como fonte de dados. Esses dados alimentam um pipeline em tempo real e um pipeline em lote. Posteriormente, as informações de ambos os pipelines são consolidadas no Hologres, que fornece uma camada de serviço unificada.
Pipeline em tempo real: O Realtime Compute for Apache Flink processa dados do Simple Log Service (SLS) em tempo real e os grava no Hologres. O Flink é um poderoso mecanismo de processamento de fluxo. O Hologres suporta gravação e atualização de dados em tempo real, permitindo consultas imediatas após a escrita. Sua integração nativa viabiliza o desenvolvimento de data warehouses em tempo real com alto throughput, baixa latência e alta qualidade, orientados por modelos. Isso atende às necessidades de insights de negócios em tempo real, como extração dos eventos mais recentes e análise de tendências.
Pipeline em lote: O MaxCompute processa e arquiva grandes volumes de dados em lote. O Object Storage Service (OSS) é um serviço do Alibaba Cloud para armazenar diversos tipos de dados. Como os dados brutos deste tutorial estão no formato JSON, o OSS oferece um armazenamento prático, seguro, econômico e confiável. O MaxCompute é um data warehouse em nuvem SaaS de nível empresarial projetado para análise de dados. Ele lê e analisa diretamente dados semiestruturados no OSS por meio de tabelas externas, integra dados de alto valor ao seu armazenamento interno e se conecta ao DataWorks para criar um data warehouse em lote.
O Hologres integra-se perfeitamente ao MaxCompute na camada de armazenamento. Isso permite usar o Hologres para acelerar consultas e análises sobre grandes volumes de dados históricos no MaxCompute, atendendo a demandas de consultas esporádicas e de alto desempenho em dados históricos. Além disso, é possível utilizar facilmente o pipeline em lote para corrigir dados em tempo real, resolvendo problemas como omissões que podem ocorrer no pipeline em tempo real.
Esta solução oferece as seguintes vantagens:
Pipeline em lote estável e eficiente: Suporta gravação e atualização horária de dados, processa grandes volumes em lote, executa cálculos e análises complexas, reduz custos computacionais e aumenta a eficiência do processamento.
Pipeline em tempo real maduro: Permite ingestão, computação e análise de eventos em tempo real. O pipeline simplificado entrega respostas em segundos.
Armazenamento e serviço unificados: O Hologres fornece uma camada de serviço unificada com armazenamento centralizado e uma interface externa consistente (uma única interface SQL para consultas OLAP e Key-Value).
Análise unificada em lote e em tempo real: Reduz redundância e movimentação de dados, além de permitir correções.
Essa abordagem de desenvolvimento integrado proporciona resposta de dados em nível de segundo, visibilidade de status ponta a ponta, arquitetura simplificada com menos componentes e dependências, além de reduções efetivas nos custos operacionais e de mão de obra.
Entendendo o negócio e os dados
Desenvolvedores geram muitos eventos ao trabalhar em projetos open source no GitHub. O GitHub registra detalhes de cada evento, incluindo tipo, desenvolvedor e repositório de código. Eventos públicos, como marcar um repositório com estrela ou fazer commit de código, ficam disponíveis publicamente. Para obter a lista completa de tipos de eventos, consulte Eventos e payloads de Webhook.
O GitHub disponibiliza eventos públicos por meio de uma OpenAPI. A API fornece dados em tempo real com atraso de cinco minutos. Para mais informações, consulte Eventos.
O projeto GH Archive coleta e disponibiliza arquivos horários de eventos públicos do GitHub. Utilize esses arquivos para obter dados offline. Para mais informações, consulte GH Archive.
Entendendo o negócio do GitHub
O negócio principal do GitHub envolve gerenciamento de código e interações, abrangendo três entidades principais: Desenvolvedor, Repositório e Organização.
Para esta análise de dados, um Event também é armazenado e registrado como entidade.

Entendendo os dados brutos de eventos públicos
O exemplo abaixo mostra os dados JSON de um evento bruto:
{
"id": "19541192931",
"type": "WatchEvent",
"actor":
{
"id": 23286640,
"login": "herekeo",
"display_login": "herekeo",
"gravatar_id": "",
"url": "https://api.github.com/users/herekeo",
"avatar_url": "https://avatars.githubusercontent.com/u/23286640?"
},
"repo":
{
"id": 52760178,
"name": "crazyguitar/pysheeet",
"url": "https://api.github.com/repos/crazyguitar/pysheeet"
},
"payload":
{
"action": "started"
},
"public": true,
"created_at": "2022-01-01T00:03:04Z"
}
Esta análise abrange 15 tipos de eventos públicos, excluindo eventos que nunca ocorreram ou não são mais registrados. Para detalhes sobre esses tipos, consulte Tipos de eventos públicos do Github.
Pré-requisitos
Uma instância Elastic Compute Service (ECS) criada e associada a um Elastic IP Address (EIP). Essa instância extrairá dados de eventos em tempo real da API do GitHub. Para mais informações, consulte Guia de criação e Elastic IP Address.
Object Storage Service (OSS) ativado e ferramenta ossutil instalada na instância ECS para armazenar arquivos de dados JSON do GH Archive. Para mais informações, consulte Ativar o OSS e Instalar o ossutil.
MaxCompute ativado e projeto criado. Para mais informações, consulte Criar um projeto do MaxCompute.
DataWorks ativado e workspace criado para configurar tarefas de agendamento offline. Para mais informações, consulte Criar um workspace.
Simple Log Service (SLS) ativado, com projeto e Logstore criados para coletar logs da instância ECS. Para mais informações, consulte Coletar e analisar logs de texto do ECS usando LoongCollector.
Instância do Realtime Compute for Apache Flink ativada para gravar dados de log do SLS no Hologres em tempo real. Para mais informações, consulte Ativar o Realtime Compute for Apache Flink.
Hologres ativado. Para mais informações, consulte Comprar uma instância do Hologres.
Construir um data warehouse offline (atualizações horárias)
Baixar arquivos de dados brutos usando uma instância ECS e enviá-los ao OSS
Utilize uma instância Elastic Compute Service (ECS) para baixar arquivos de dados JSON do GH Archive.
Baixe dados históricos com o comando
wget. Por exemplo, executewget https://data.gharchive.org/{2012..2022}-{01..12}-{01..31}-{0..23}.json.gzpara baixar dados horários de 2012 a 2022.-
Para baixar novos dados gerados a cada hora, configure uma tarefa agendada conforme descrito abaixo.
NotaCertifique-se de que o ossutil esteja instalado na instância ECS. Para mais informações, consulte Instalar o ossutil. Baixe o pacote de instalação do ossutil e faça upload dele para a instância ECS. Execute
yum install unzippara instalar o software unzip. Em seguida, descompacte o pacote do ossutil e mova o arquivo executável para o diretório/usr/bin/.Verifique se você criou um bucket do Object Storage Service (OSS) na mesma região da sua instância ECS. É possível usar um nome personalizado para o bucket. Este exemplo utiliza o nome
githubevents.Neste exemplo, os arquivos são baixados para o diretório
/opt/hourlydata/gh_datana instância ECS. Você pode escolher outro diretório se preferir.
-
Execute o comando a seguir para criar um arquivo chamado
download_code.shno diretório/opt/hourlydata.cd /opt/hourlydata vim download_code.sh -
Pressione
ipara entrar no modo de edição e adicione o script a seguir.d=$(TZ=UTC date --date='1 hour ago' '+%Y-%m-%d-%-H') h=$(TZ=UTC date --date='1 hour ago' '+%Y-%m-%d-%H') url=https://data.gharchive.org/${d}.json.gz echo ${url} # Download the data to the ./gh_data/ directory. You can use a different directory. wget ${url} -P ./gh_data/ # Change to the gh_data directory. cd gh_data # Decompress the downloaded data into a JSON file. gzip -d ${d}.json echo ${d}.json # Change to the root directory. cd /root # Use ossutil to upload the data to OSS. # Create the hr=${h} directory in the githubevents OSS bucket. ossutil mkdir oss://githubevents/hr=${h} # Upload the data from the /opt/hourlydata/gh_data directory to OSS. You can use a different directory. ossutil cp -r /opt/hourlydata/gh_data oss://githubevents/hr=${h} -u echo oss uploaded successfully! rm -rf /opt/hourlydata/gh_data/${d}.json echo ecs deleted! Pressione a tecla Esc, insira
:wqe pressione Enter para salvar e fechar o arquivo.-
Execute o comando a seguir para rodar o script
download_code.shaos 10 minutos de cada hora.# 1. Run the following command and press I to enter edit mode. crontab -e # 2. Add the following command. Then, press Esc, enter :wq, and press Enter to exit. 10 * * * * cd /opt/hourlydata && sh download_code.sh > download.logApós a execução do script, o arquivo JSON da hora anterior é baixado aos 10 minutos de cada hora. O arquivo é então descompactado na instância ECS e enviado ao OSS no caminho
oss://githubevents. Para ler apenas o arquivo da hora anterior, cria-se um diretório chamado'hr=%Y-%M-%D-%H'como partição para cada arquivo durante o upload. Isso garante que as operações subsequentes de gravação leiam arquivos apenas da partição mais recente.
Importar dados do OSS para o MaxCompute usando tabela externa
Execute os comandos a seguir no cliente MaxCompute ou em um nó ODPS SQL no DataWorks. Para mais informações, consulte Conectar ao MaxCompute usando o cliente (odpscmd) ou Desenvolver uma tarefa ODPS SQL.
-
Crie a tabela externa
githubeventspara ler os arquivos JSON armazenados no OSS:CREATE EXTERNAL TABLE IF NOT EXISTS githubevents ( col STRING ) PARTITIONED BY ( hr STRING ) STORED AS textfile LOCATION 'oss://oss-cn-hangzhou-internal.aliyuncs.com/githubevents/' ;Para mais informações sobre como criar tabelas externas para acessar dados do OSS no MaxCompute, consulte Acessar dados não estruturados no OSS.
-
Crie a tabela de fatos
dwd_github_events_odpspara armazenar os dados. O código a seguir mostra a instrução DDL (Data Definition Language):CREATE TABLE IF NOT EXISTS dwd_github_events_odps ( id BIGINT COMMENT 'Event ID' ,actor_id BIGINT COMMENT 'ID of the event initiator' ,actor_login STRING COMMENT 'Logon name of the event initiator' ,repo_id BIGINT COMMENT 'Repository ID' ,repo_name STRING COMMENT 'Full name of the repository in the format of owner/repository_name' ,org_id BIGINT COMMENT 'ID of the organization to which the repository belongs' ,org_login STRING COMMENT 'Name of the organization to which the repository belongs' ,`type` STRING COMMENT 'Event type' ,created_at DATETIME COMMENT 'Time when the event occurred' ,action STRING COMMENT 'Event action' ,iss_or_pr_id BIGINT COMMENT 'ID of the issue or pull request' ,number BIGINT COMMENT 'Number of the issue or pull request' ,comment_id BIGINT COMMENT 'Comment ID' ,commit_id STRING COMMENT 'Commit ID' ,member_id BIGINT COMMENT 'Member ID' ,rev_or_push_or_rel_id BIGINT COMMENT 'ID of the review, push, or release' ,ref STRING COMMENT 'Name of the created or deleted resource' ,ref_type STRING COMMENT 'Type of the created or deleted resource' ,state STRING COMMENT 'Status of the issue, pull request, or pull request review' ,author_association STRING COMMENT 'Relationship between the actor and the repository' ,language STRING COMMENT 'Language of the code in the pull request' ,merged BOOLEAN COMMENT 'Indicates whether the pull request was merged' ,merged_at DATETIME COMMENT 'Time when the code was merged' ,additions BIGINT COMMENT 'Number of added lines of code' ,deletions BIGINT COMMENT 'Number of deleted lines of code' ,changed_files BIGINT COMMENT 'Number of files changed in the pull request' ,push_size BIGINT COMMENT 'Number of commits' ,push_distinct_size BIGINT COMMENT 'Number of distinct commits' ,hr STRING COMMENT 'Hour when the event occurred. For example, if the event occurred at 00:23, the value of hr is 00.' ,`month` STRING COMMENT 'Month when the event occurred. For example, if the event occurred in October 2015, the value of month is 2015-10.' ,`year` STRING COMMENT 'Year when the event occurred. For example, if the event occurred in 2015, the value of year is 2015.' ) PARTITIONED BY ( ds STRING COMMENT 'Date when the event occurred, in the yyyy-mm-dd format.' ); -
Analise os dados JSON e grave-os na tabela de fatos.
Execute o comando a seguir para adicionar partições, analisar os dados JSON e gravar os dados na tabela
dwd_github_events_odps:msck repair table githubevents add partitions; set odps.sql.hive.compatible = true; set odps.sql.split.hive.bridge = true; INSERT into TABLE dwd_github_events_odps PARTITION(ds) SELECT CAST(GET_JSON_OBJECT(col,'$.id') AS BIGINT ) AS id ,CAST(GET_JSON_OBJECT(col,'$.actor.id')AS BIGINT) AS actor_id ,GET_JSON_OBJECT(col,'$.actor.login') AS actor_login ,CAST(GET_JSON_OBJECT(col,'$.repo.id')AS BIGINT) AS repo_id ,GET_JSON_OBJECT(col,'$.repo.name') AS repo_name ,CAST(GET_JSON_OBJECT(col,'$.org.id')AS BIGINT) AS org_id ,GET_JSON_OBJECT(col,'$.org.login') AS org_login ,GET_JSON_OBJECT(col,'$.type') as type ,to_date(GET_JSON_OBJECT(col,'$.created_at'), 'yyyy-mm-ddThh:mi:ssZ') AS created_at ,GET_JSON_OBJECT(col,'$.payload.action') AS action ,case WHEN GET_JSON_OBJECT(col,'$.type')="PullRequestEvent" THEN CAST(GET_JSON_OBJECT(col,'$.payload.pull_request.id')AS BIGINT) WHEN GET_JSON_OBJECT(col,'$.type')="IssuesEvent" THEN CAST(GET_JSON_OBJECT(col,'$.payload.issue.id')AS BIGINT) END AS iss_or_pr_id ,case WHEN GET_JSON_OBJECT(col,'$.type')="PullRequestEvent" THEN CAST(GET_JSON_OBJECT(col,'$.payload.pull_request.number')AS BIGINT) WHEN GET_JSON_OBJECT(col,'$.type')="IssuesEvent" THEN CAST(GET_JSON_OBJECT(col,'$.payload.issue.number')AS BIGINT) ELSE CAST(GET_JSON_OBJECT(col,'$.payload.number')AS BIGINT) END AS number ,CAST(GET_JSON_OBJECT(col,'$.payload.comment.id')AS BIGINT) AS comment_id ,GET_JSON_OBJECT(col,'$.payload.comment.commit_id') AS commit_id ,CAST(GET_JSON_OBJECT(col,'$.payload.member.id')AS BIGINT) AS member_id ,case WHEN GET_JSON_OBJECT(col,'$.type')="PullRequestReviewEvent" THEN CAST(GET_JSON_OBJECT(col,'$.payload.review.id')AS BIGINT) WHEN GET_JSON_OBJECT(col,'$.type')="PushEvent" THEN CAST(GET_JSON_OBJECT(col,'$.payload.push_id')AS BIGINT) WHEN GET_JSON_OBJECT(col,'$.type')="ReleaseEvent" THEN CAST(GET_JSON_OBJECT(col,'$.payload.release.id')AS BIGINT) END AS rev_or_push_or_rel_id ,GET_JSON_OBJECT(col,'$.payload.ref') AS ref ,GET_JSON_OBJECT(col,'$.payload.ref_type') AS ref_type ,case WHEN GET_JSON_OBJECT(col,'$.type')="PullRequestEvent" THEN GET_JSON_OBJECT(col,'$.payload.pull_request.state') WHEN GET_JSON_OBJECT(col,'$.type')="IssuesEvent" THEN GET_JSON_OBJECT(col,'$.payload.issue.state') WHEN GET_JSON_OBJECT(col,'$.type')="PullRequestReviewEvent" THEN GET_JSON_OBJECT(col,'$.payload.review.state') END AS state ,case WHEN GET_JSON_OBJECT(col,'$.type')="PullRequestEvent" THEN GET_JSON_OBJECT(col,'$.payload.pull_request.author_association') WHEN GET_JSON_OBJECT(col,'$.type')="IssuesEvent" THEN GET_JSON_OBJECT(col,'$.payload.issue.author_association') WHEN GET_JSON_OBJECT(col,'$.type')="IssueCommentEvent" THEN GET_JSON_OBJECT(col,'$.payload.comment.author_association') WHEN GET_JSON_OBJECT(col,'$.type')="PullRequestReviewEvent" THEN GET_JSON_OBJECT(col,'$.payload.review.author_association') END AS author_association ,GET_JSON_OBJECT(col,'$.payload.pull_request.base.repo.language') AS language ,CAST(GET_JSON_OBJECT(col,'$.payload.pull_request.merged') AS BOOLEAN) AS merged ,to_date(GET_JSON_OBJECT(col,'$.payload.pull_request.merged_at'), 'yyyy-mm-ddThh:mi:ssZ') AS merged_at ,CAST(GET_JSON_OBJECT(col,'$.payload.pull_request.additions')AS BIGINT) AS additions ,CAST(GET_JSON_OBJECT(col,'$.payload.pull_request.deletions')AS BIGINT) AS deletions ,CAST(GET_JSON_OBJECT(col,'$.payload.pull_request.changed_files')AS BIGINT) AS changed_files ,CAST(GET_JSON_OBJECT(col,'$.payload.size')AS BIGINT) AS push_size ,CAST(GET_JSON_OBJECT(col,'$.payload.distinct_size')AS BIGINT) AS push_distinct_size ,SUBSTR(GET_JSON_OBJECT(col,'$.created_at'),12,2) as hr ,REPLACE(SUBSTR(GET_JSON_OBJECT(col,'$.created_at'),1,7),'/','-') as month ,SUBSTR(GET_JSON_OBJECT(col,'$.created_at'),1,4) as year ,REPLACE(SUBSTR(GET_JSON_OBJECT(col,'$.created_at'),1,10),'/','-') as ds from githubevents where hr = cast(to_char(dateadd(getdate(),-9,'hh'), 'yyyy-mm-dd-hh') as string); -
Consulte os dados.
Execute o comando a seguir para consultar dados da tabela
dwd_github_events_odps:SET odps.sql.allow.fullscan=true; SELECT * FROM dwd_github_events_odps where ds = '2023-03-31' limit 10;O resultado de exemplo a seguir é retornado:

Construir um data warehouse em tempo real
Obter dados em tempo real usando ECS
Uma instância Elastic Compute Service (ECS) é usada para extrair dados de eventos em tempo real da API do GitHub. Este tópico fornece um script de exemplo para demonstrar como coletar dados em tempo real da API do GitHub.
Cada vez que o script é executado, ele roda por 1 minuto. Durante esse período, ele coleta os dados de eventos em tempo real fornecidos pela API e armazena cada evento no formato JSON.
Este script não garante a coleta de todos os dados de eventos em tempo real.
Para coletar dados continuamente da API do GitHub, forneça os cabeçalhos Accept e Authorization. O valor de Accept é fixo. Para Authorization, insira o token de acesso pessoal obtido no GitHub. Para mais informações sobre como criar um token de acesso pessoal, consulte esta documentação.
-
Execute os comandos a seguir para criar um arquivo chamado
download_realtime_data.pyno diretório/opt/realtime.cd /opt/realtime vim download_realtime_data.py -
Pressione
ipara entrar no modo de edição e adicione o conteúdo de exemplo a seguir ao arquivo.#!python import requests import json import sys import time # Get the API URL def get_next_link(resp): resp_link = resp.headers['link'] link = '' for l in resp_link.split(', '): link = l.split('; ')[0][1:-1] rel = l.split('; ')[1] if rel == 'rel="next"': return link return None # Collect one page of data from the API def download(link, fname): # Define the Accept and Authorization headers for the GitHub API headers = {"Accept": "application/vnd.github+json","Authorization": "<Bearer> <github_api_token>"} resp = requests.get(link, headers=headers) if int(resp.status_code) != 200: return None with open(fname, 'a') as f: for j in resp.json(): f.write(json.dumps(j)) f.write('\n') print('downloaded {} events to {}'.format(len(resp.json()), fname)) return resp # Collect multiple pages of data from the API def download_all_data(fname): link = 'https://api.github.com/events?per_page=100&page=1' while True: resp = download(link, fname) if resp is None: break link = get_next_link(resp) if link is None: break # Define the current time def get_current_ms(): return round(time.time()*1000) # Define the script execution duration as 1 minute def main(fname): current_ms = get_current_ms() while get_current_ms() - current_ms < 60*1000: download_all_data(fname) time.sleep(0.1) # Run the script if __name__ == '__main__': if len(sys.argv) < 2: print('usage: python {} <log_file>'.format(sys.argv[0])) exit(0) main(sys.argv[1]) -
Pressione a tecla Esc, insira
:wqe pressione Enter para salvar e fechar o arquivo. -
Crie um arquivo
run_py.shpara executar odownload_realtime_data.pye armazenar separadamente os dados coletados em cada execução. O conteúdo é o seguinte:python /opt/realtime/download_realtime_data.py /opt/realtime/gh_realtime_data/$(date '+%Y-%m-%d-%H:%M:%S').json -
Crie um arquivo
delete_log.shpara excluir dados históricos. O conteúdo é o seguinte:d=$(TZ=UTC date --date='2 day ago' '+%Y-%m-%d') rm -f /opt/realtime/gh_realtime_data/*${d}*.json -
Execute os comandos a seguir para coletar dados do GitHub a cada minuto e excluir dados históricos diariamente.
#1. Run the following command and press I to enter edit mode. crontab -e #2. Add the following commands. Then, press Esc, enter :wq, and press Enter to exit. * * * * * bash /opt/realtime/run_py.sh 1 1 * * * bash /opt/realtime/delete_log.sh
Coletar dados do ECS usando SLS
O Simple Log Service (SLS) é utilizado para coletar os dados de eventos em tempo real extraídos da instância ECS como logs.
O SLS suporta a coleta de logs de instâncias ECS por meio do Logtail. Como os dados neste tópico estão no formato JSON, utilize o modo JSON do Logtail para coletar rapidamente logs JSON incrementais da instância ECS. Para mais informações, consulte Coletar logs no modo JSON. Neste tópico, o SLS é configurado para analisar os pares chave-valor de nível superior dos dados brutos.
Neste exemplo, o parâmetro de caminho de log da configuração do Logtail está definido como /opt/realtime/gh_realtime_data/**/*.json.
Após concluir a configuração, o SLS coleta continuamente dados incrementais de eventos da instância ECS. A figura a seguir mostra um exemplo dos dados coletados.
Gravar dados do SLS no Hologres em tempo real usando Flink
O Flink é usado para gravar dados de log coletados pelo SLS no Hologres em tempo real. Ao utilizar uma tabela de origem SLS e uma tabela de resultados Hologres no Flink, é possível gravar dados do SLS no Hologres em tempo real. Para mais informações, consulte Importar dados do Simple Log Service.
-
Crie uma tabela interna do Hologres.
A tabela interna criada neste tópico retém apenas alguns pares chave-valor dos dados JSON brutos. O
iddo evento e a datadssão definidos como chave primária. Oiddo evento é definido como chave de distribuição. A datadsé definida como chave de partição. O tempo do eventocreated_até definido como event_time_column. Crie índices para outros campos conforme necessário para melhorar a eficiência das consultas. Para mais informações sobre índices, consulte CREATE TABLE. A seguinte instrução DDL (Data Definition Language) é usada para criar a tabela neste exemplo.DROP TABLE IF EXISTS gh_realtime_data; BEGIN; CREATE TABLE gh_realtime_data ( id bigint, actor_id bigint, actor_login text, repo_id bigint, repo_name text, org_id bigint, org_login text, type text, created_at timestamp with time zone NOT NULL, action text, iss_or_pr_id bigint, number bigint, comment_id bigint, commit_id text, member_id bigint, rev_or_push_or_rel_id bigint, ref text, ref_type text, state text, author_association text, language text, merged boolean, merged_at timestamp with time zone, additions bigint, deletions bigint, changed_files bigint, push_size bigint, push_distinct_size bigint, hr text, month text, year text, ds text, PRIMARY KEY (id,ds) ) PARTITION BY LIST (ds); CALL set_table_property('public.gh_realtime_data', 'distribution_key', 'id'); CALL set_table_property('public.gh_realtime_data', 'event_time_column', 'created_at'); CALL set_table_property('public.gh_realtime_data', 'clustering_key', 'created_at'); COMMENT ON COLUMN public.gh_realtime_data.id IS 'Event ID'; COMMENT ON COLUMN public.gh_realtime_data.actor_id IS 'ID of the event initiator'; COMMENT ON COLUMN public.gh_realtime_data.actor_login IS 'Logon name of the event initiator'; COMMENT ON COLUMN public.gh_realtime_data.repo_id IS 'repo ID'; COMMENT ON COLUMN public.gh_realtime_data.repo_name IS 'repo name'; COMMENT ON COLUMN public.gh_realtime_data.org_id IS 'ID of the organization to which the repo belongs'; COMMENT ON COLUMN public.gh_realtime_data.org_login IS 'Name of the organization to which the repo belongs'; COMMENT ON COLUMN public.gh_realtime_data.type IS 'Event type'; COMMENT ON COLUMN public.gh_realtime_data.created_at IS 'Time when the event occurred'; COMMENT ON COLUMN public.gh_realtime_data.action IS 'Event action'; COMMENT ON COLUMN public.gh_realtime_data.iss_or_pr_id IS 'issue/pull_request ID'; COMMENT ON COLUMN public.gh_realtime_data.number IS 'issue/pull_request number'; COMMENT ON COLUMN public.gh_realtime_data.comment_id IS 'comment ID'; COMMENT ON COLUMN public.gh_realtime_data.commit_id IS 'Commit ID'; COMMENT ON COLUMN public.gh_realtime_data.member_id IS 'Member ID'; COMMENT ON COLUMN public.gh_realtime_data.rev_or_push_or_rel_id IS 'review/push/release ID'; COMMENT ON COLUMN public.gh_realtime_data.ref IS 'Name of the created or deleted resource'; COMMENT ON COLUMN public.gh_realtime_data.ref_type IS 'Type of the created or deleted resource'; COMMENT ON COLUMN public.gh_realtime_data.state IS 'Status of the issue/pull_request/pull_request_review'; COMMENT ON COLUMN public.gh_realtime_data.author_association IS 'Relationship between the actor and the repo'; COMMENT ON COLUMN public.gh_realtime_data.language IS 'Programming language'; COMMENT ON COLUMN public.gh_realtime_data.merged IS 'Specifies whether the merge is accepted'; COMMENT ON COLUMN public.gh_realtime_data.merged_at IS 'Time when the code was merged'; COMMENT ON COLUMN public.gh_realtime_data.additions IS 'Number of added lines of code'; COMMENT ON COLUMN public.gh_realtime_data.deletions IS 'Number of deleted lines of code'; COMMENT ON COLUMN public.gh_realtime_data.changed_files IS 'Number of files changed in the pull request'; COMMENT ON COLUMN public.gh_realtime_data.push_size IS 'Number of pushes'; COMMENT ON COLUMN public.gh_realtime_data.push_distinct_size IS 'Number of distinct pushes'; COMMENT ON COLUMN public.gh_realtime_data.hr IS 'The hour when the event occurred. For example, if the time is 00:23, hr=00.'; COMMENT ON COLUMN public.gh_realtime_data.month IS 'The month when the event occurred. For example, if the date is October 2015, month=2015-10.'; COMMENT ON COLUMN public.gh_realtime_data.year IS 'The year when the event occurred. For example, if the year is 2015, year=2015.'; COMMENT ON COLUMN public.gh_realtime_data.ds IS 'The day when the event occurred. ds=yyyy-mm-dd.'; COMMIT; -
Grave dados em tempo real usando o Flink.
Use o Flink para analisar ainda mais os dados do SLS e gravá-los no Hologres em tempo real. Utilize as instruções a seguir no Flink para filtrar os dados a serem gravados. Isso descarta dados incorretos onde o ID do evento ou o tempo do evento (
created_at) é nulo e retém apenas dados de eventos recentes.CREATE TEMPORARY TABLE sls_input ( actor varchar, created_at varchar, id bigint, org varchar, payload varchar, public varchar, repo varchar, type varchar ) WITH ( 'connector' = 'sls', 'endpoint' = '<endpoint>',--The private endpoint of SLS 'accessid' = '<accesskey id>',--The AccessKey ID of your account 'accesskey' = '<accesskey secret>',--The AccessKey secret of your account 'project' = '<project name>',--The name of the SLS project 'logstore' = '<logstore name>'--The name of the SLS Logstore 'starttime' = '2023-04-06 00:00:00',--The start time for SLS data collection ); CREATE TEMPORARY TABLE hologres_sink ( id bigint, actor_id bigint, actor_login string, repo_id bigint, repo_name string, org_id bigint, org_login string, type string, created_at timestamp, action string, iss_or_pr_id bigint, number bigint, comment_id bigint, commit_id string, member_id bigint, rev_or_push_or_rel_id bigint, `ref` string, ref_type string, state string, author_association string, `language` string, merged boolean, merged_at timestamp, additions bigint, deletions bigint, changed_files bigint, push_size bigint, push_distinct_size bigint, hr string, `month` string, `year` string, ds string ) WITH ( 'connector' = 'hologres', 'dbname' = '<hologres dbname>', --The name of the Hologres database 'tablename' = '<hologres tablename>', --The name of the Hologres table that receives data 'username' = '<accesskey id>', --The AccessKey ID of the current Alibaba Cloud account 'password' = '<accesskey secret>', --The AccessKey Secret of the current Alibaba Cloud account 'endpoint' = '<endpoint>', --The VPC endpoint of the current Hologres instance 'jdbcretrycount' = '1', --The number of retries upon connection failure 'partitionrouter' = 'true', --Specifies whether to write data to a partitioned table 'createparttable' = 'true', --Specifies whether to automatically create partitions 'mutatetype' = 'insertorignore' --The data writing mode ); INSERT INTO hologres_sink SELECT id ,CAST(JSON_VALUE(actor, '$.id') AS bigint) AS actor_id ,JSON_VALUE(actor, '$.login') AS actor_login ,CAST(JSON_VALUE(repo, '$.id') AS bigint) AS repo_id ,JSON_VALUE(repo, '$.name') AS repo_name ,CAST(JSON_VALUE(org, '$.id') AS bigint) AS org_id ,JSON_VALUE(org, '$.login') AS org_login ,type ,TO_TIMESTAMP_TZ(replace(created_at,'T',' '), 'yyyy-MM-dd HH:mm:ss', 'UTC') AS created_at ,JSON_VALUE(payload, '$.action') AS action ,CASE WHEN type='PullRequestEvent' THEN CAST(JSON_VALUE(payload, '$.pull_request.id') AS bigint) WHEN type='IssuesEvent' THEN CAST(JSON_VALUE(payload, '$.issue.id') AS bigint) END AS iss_or_pr_id ,CASE WHEN type='PullRequestEvent' THEN CAST(JSON_VALUE(payload, '$.pull_request.number') AS bigint) WHEN type='IssuesEvent' THEN CAST(JSON_VALUE(payload, '$.issue.number') AS bigint) ELSE CAST(JSON_VALUE(payload, '$.number') AS bigint) END AS number ,CAST(JSON_VALUE(payload, '$.comment.id') AS bigint) AS comment_id ,JSON_VALUE(payload, '$.comment.commit_id') AS commit_id ,CAST(JSON_VALUE(payload, '$.member.id') AS bigint) AS member_id ,CASE WHEN type='PullRequestReviewEvent' THEN CAST(JSON_VALUE(payload, '$.review.id') AS bigint) WHEN type='PushEvent' THEN CAST(JSON_VALUE(payload, '$.push_id') AS bigint) WHEN type='ReleaseEvent' THEN CAST(JSON_VALUE(payload, '$.release.id') AS bigint) END AS rev_or_push_or_rel_id ,JSON_VALUE(payload, '$.ref') AS `ref` ,JSON_VALUE(payload, '$.ref_type') AS ref_type ,CASE WHEN type='PullRequestEvent' THEN JSON_VALUE(payload, '$.pull_request.state') WHEN type='IssuesEvent' THEN JSON_VALUE(payload, '$.issue.state') WHEN type='PullRequestReviewEvent' THEN JSON_VALUE(payload, '$.review.state') END AS state ,CASE WHEN type='PullRequestEvent' THEN JSON_VALUE(payload, '$.pull_request.author_association') WHEN type='IssuesEvent' THEN JSON_VALUE(payload, '$.issue.author_association') WHEN type='IssueCommentEvent' THEN JSON_VALUE(payload, '$.comment.author_association') WHEN type='PullRequestReviewEvent' THEN JSON_VALUE(payload, '$.review.author_association') END AS author_association ,JSON_VALUE(payload, '$.pull_request.base.repo.language') AS `language` ,CAST(JSON_VALUE(payload, '$.pull_request.merged') AS boolean) AS merged ,TO_TIMESTAMP_TZ(replace(JSON_VALUE(payload, '$.pull_request.merged_at'),'T',' '), 'yyyy-MM-dd HH:mm:ss', 'UTC') AS merged_at ,CAST(JSON_VALUE(payload, '$.pull_request.additions') AS bigint) AS additions ,CAST(JSON_VALUE(payload, '$.pull_request.deletions') AS bigint) AS deletions ,CAST(JSON_VALUE(payload, '$.pull_request.changed_files') AS bigint) AS changed_files ,CAST(JSON_VALUE(payload, '$.size') AS bigint) AS push_size ,CAST(JSON_VALUE(payload, '$.distinct_size') AS bigint) AS push_distinct_size ,SUBSTRING(TO_TIMESTAMP_TZ(replace(created_at,'T',' '), 'yyyy-MM-dd HH:mm:ss', 'UTC'),12,2) as hr ,REPLACE(SUBSTRING(TO_TIMESTAMP_TZ(replace(created_at,'T',' '), 'yyyy-MM-dd HH:mm:ss', 'UTC'),1,7),'/','-') as `month` ,SUBSTRING(TO_TIMESTAMP_TZ(replace(created_at,'T',' '), 'yyyy-MM-dd HH:mm:ss', 'UTC'),1,4) as `year` ,SUBSTRING(TO_TIMESTAMP_TZ(replace(created_at,'T',' '), 'yyyy-MM-dd HH:mm:ss', 'UTC'),1,10) as ds FROM sls_input WHERE id IS NOT NULL AND created_at IS NOT NULL AND to_date(replace(created_at,'T',' ')) >= date_add(CURRENT_DATE, -1);Para mais informações sobre os parâmetros, consulte Simple Log Service (SLS) e Hologres.
NotaOs dados brutos de eventos do GitHub usam o fuso horário UTC e não possuem atributo de fuso horário. O fuso horário padrão do Hologres é UTC+8. Portanto, ajuste o fuso horário ao gravar dados do Flink para o Hologres em tempo real. Atribua o atributo de fuso horário UTC aos dados da tabela de origem no Flink SQL. Os passos são os seguintes:
Passo 1: Acessar a página de edição de jobs
Faça login no Realtime Compute for Apache Flink console
Acesse o workspace desejado
Encontre seu job Flink SQL ou JAR e clique em Edit
Passo 2: Abrir a aba Deployment Details
⚠️ Mudança importante: "Flink Configuration" não é mais uma seção separada. Ela foi mesclada em "Deployment Details".
No topo da página de edição do job, mude para a aba Deployment Details.
Role para baixo até a seção Parameter Configuration.
Passo 3: Adicionar uma configuração personalizada
Clique no botão Edit à direita de Parameter Configuration.
Na caixa de diálogo exibida, localize a caixa de texto Other Configuration.
Na caixa de texto, adicione o parâmetro do Flink
table.local-time-zone:Asia/Shanghaicomo um par chave-valor para definir o fuso horário do sistema Flink comoAsia/Shanghai.
-
Consulte os dados.
Consulte os dados do SLS gravados no Hologres através do Flink. Em seguida, realize o desenvolvimento de dados conforme necessário.
SELECT * FROM public.gh_realtime_data limit 10;O resultado a seguir é um exemplo:

Corrigir dados em tempo real usando dados offline
No cenário descrito neste tópico, dados em tempo real podem estar ausentes. Utilize dados offline para corrigir os dados em tempo real. As etapas a seguir mostram como corrigir os dados em tempo real do dia anterior. Ajuste o período de correção de dados conforme necessário.
-
Crie uma tabela estrangeira no Hologres para obter dados offline do MaxCompute.
IMPORT FOREIGN SCHEMA <maxcompute_project_name> LIMIT to ( <foreign_table_name> ) FROM SERVER odps_server INTO public OPTIONS(if_table_exist 'update',if_unsupported_type 'error');Para mais informações sobre os parâmetros, consulte IMPORT FOREIGN SCHEMA.
-
Crie uma tabela temporária para corrigir os dados em tempo real do dia anterior com dados offline.
NotaO Hologres V2.1.17 e versões posteriores suportam Serverless Computing. Para cenários como importação de dados offline em larga escala, grandes jobs ETL e consultas de alto volume em tabelas estrangeiras, utilize o Serverless Computing para executar essas tarefas. Esse recurso usa recursos serverless adicionais em vez dos recursos da sua instância, o que melhora a estabilidade da instância e reduz a probabilidade de erros de falta de memória (OOM). Não é necessário reservar recursos computacionais extras para sua instância, e você paga apenas pelas tarefas executadas. Para mais informações sobre Serverless Computing, consulte Serverless Computing. Para instruções sobre como usar o Serverless Computing, consulte Usar Serverless Computing.
-- Clean up potential temporary tables DROP TABLE IF EXISTS gh_realtime_data_tmp; -- Create a temporary table SET hg_experimental_enable_create_table_like_properties = ON; CALL HG_CREATE_TABLE_LIKE ('gh_realtime_data_tmp', 'select * from gh_realtime_data'); -- (Optional) Use Serverless Computing to perform large-scale offline data import and ETL jobs. SET hg_computing_resource = 'serverless'; -- Insert data into the temporary table and update statistics INSERT INTO gh_realtime_data_tmp SELECT * FROM <foreign_table_name> WHERE ds = current_date - interval '1 day' ON CONFLICT (id, ds) DO NOTHING; ANALYZE gh_realtime_data_tmp; -- Reset the configuration to ensure that non-essential SQL statements do not use Serverless resources. RESET hg_computing_resource; -- Replace the atomic table with the existing temporary child table BEGIN; DROP TABLE IF EXISTS "gh_realtime_data_<yesterday_date>"; ALTER TABLE gh_realtime_data_tmp RENAME TO "gh_realtime_data_<yesterday_date>"; ALTER TABLE gh_realtime_data ATTACH PARTITION "gh_realtime_data_<yesterday_date>" FOR VALUES IN ('<yesterday_date>'); COMMIT;
Análise de dados
É possível realizar uma ampla variedade de análises sobre o grande volume de dados coletados. Com base no intervalo de tempo exigido pelo seu negócio, projete seu data warehouse em camadas. Isso suporta diversas necessidades, como análise em tempo real, análise offline e análise integrada em tempo real e offline.
Os exemplos a seguir analisam os dados em tempo real obtidos anteriormente. Também é possível analisar dados de repositórios de código ou desenvolvedores específicos.
-
Consulte o número total de eventos públicos de hoje.
SELECT count(*) FROM gh_realtime_data WHERE created_at >= date_trunc('day', now());Segue um resultado de exemplo:
count ------ 1006 -
Identifique os projetos mais ativos (com mais eventos) no último dia.
SELECT repo_name, COUNT(*) AS events FROM gh_realtime_data WHERE created_at >= now() - interval '1 day' GROUP BY repo_name ORDER BY events DESC LIMIT 5;Segue um resultado de exemplo:
repo_name events ----------------------------------------+------ leo424y/heysiri.ml 29 arm-on/plan 10 Christoffel-T/fiverr-pat-20230331 9 mate-academy/react_dynamic-list-of-goods 9 openvinotoolkit/openvino 7 -
Liste os desenvolvedores mais ativos (com mais eventos) no último dia.
SELECT actor_login, COUNT(*) AS events FROM gh_realtime_data WHERE created_at >= now() - interval '1 day' AND actor_login NOT LIKE '%[bot]' GROUP BY actor_login ORDER BY events DESC LIMIT 5;Segue um resultado de exemplo:
actor_login events ------------------+------ direwolf-github 13 arm-on 10 sergii-nosachenko 9 Christoffel-T 9 yangwang201911 7 -
Verifique a classificação das linguagens de programação mais populares na última hora.
SELECT language, count(*) total FROM gh_realtime_data WHERE created_at > now() - interval '1 hour' AND language IS NOT NULL GROUP BY language ORDER BY total DESC LIMIT 10;Segue um resultado de exemplo:
language total -----------+---- JavaScript 25 C++ 15 Python 14 TypeScript 13 Java 8 PHP 8 -
Classifique os projetos pelo número de estrelas recebidas no último dia.
NotaEste exemplo não considera casos em que usuários removem a estrela de um projeto.
SELECT repo_id, repo_name, COUNT(actor_login) total FROM gh_realtime_data WHERE type = 'WatchEvent' AND created_at > now() - interval '1 day' GROUP BY repo_id, repo_name ORDER BY total DESC LIMIT 10;Segue um resultado de exemplo:
repo_id repo_name total ---------+----------------------------------+----- 618058471 facebookresearch/segment-anything 4 619959033 nomic-ai/gpt4all 1 97249406 denysdovhan/wtfjs 1 9791525 digininja/DVWA 1 168118422 aylei/interview 1 343520006 joehillen/sysz 1 162279822 agalwood/Motrix 1 577723410 huggingface/swift-coreml-diffusers 1 609539715 e2b-dev/e2b 1 254839429 maniackk/KKCallStack 1 -
Consulte os usuários ativos diários e projetos de hoje.
SELECT uniq (actor_id) actor_num, uniq (repo_id) repo_num FROM gh_realtime_data WHERE created_at > date_trunc('day', now());Segue um resultado de exemplo:
actor_num repo_num ---------+-------- 743 816