Todos os produtos
Search
Central de documentação

MaxCompute:Implementando análise unificada em lote e em tempo real usando o conjunto de dados público de eventos do GitHub

Última atualização: Jul 04, 2026

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.

image

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

Para esta análise de dados, um Event também é armazenado e registrado como entidade.

image

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

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, execute wget https://data.gharchive.org/{2012..2022}-{01..12}-{01..31}-{0..23}.json.gz para baixar dados horários de 2012 a 2022.

  • Para baixar novos dados gerados a cada hora, configure uma tarefa agendada conforme descrito abaixo.

    Nota
    • Certifique-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 unzip para 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_data na instância ECS. Você pode escolher outro diretório se preferir.

    1. Execute o comando a seguir para criar um arquivo chamado download_code.sh no diretório /opt/hourlydata.

      cd /opt/hourlydata
      vim download_code.sh
    2. Pressione i para 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!
    3. Pressione a tecla Esc, insira :wq e pressione Enter para salvar e fechar o arquivo.

    4. Execute o comando a seguir para rodar o script download_code.sh aos 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.log

      Apó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.

  1. Crie a tabela externa githubevents para 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.

  2. Crie a tabela de fatos dwd_github_events_odps para 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.'
    );
  3. 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);
  4. 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:

    image

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.

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

  1. Execute os comandos a seguir para criar um arquivo chamado download_realtime_data.py no diretório /opt/realtime.

    cd /opt/realtime
    vim download_realtime_data.py
  2. Pressione i para 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])
  3. Pressione a tecla Esc, insira :wq e pressione Enter para salvar e fechar o arquivo.

  4. Crie um arquivo run_py.sh para executar o download_realtime_data.py e 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
  5. Crie um arquivo delete_log.sh para 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
  6. 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.

Nota

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

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.

  1. 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 id do evento e a data ds são definidos como chave primária. O id do evento é definido como chave de distribuição. A data ds é definida como chave de partição. O tempo do evento created_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;
  2. 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.

    Nota

    Os 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/Shanghai como um par chave-valor para definir o fuso horário do sistema Flink como Asia/Shanghai.

  3. 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:

    image

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.

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

  2. Crie uma tabela temporária para corrigir os dados em tempo real do dia anterior com dados offline.

    Nota

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

    Nota

    Este 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