GitHub イベントデータを使用して、バッチとリアルタイムの統合分析ソリューションを構築します。MaxCompute はバッチデータウェアハウスとして機能し、Realtime Compute for Apache Flink と Hologres はリアルタイムデータウェアハウスを形成します。そして、Hologres と MaxCompute は、リアルタイムとバッチ両方のデータ分析のための統合レイヤーを提供します。
背景情報
ビジネスのデジタル化が進むにつれて、より新しいデータへの需要が高まっています。大規模データに対する従来のバッチ処理に加え、多くのビジネスではリアルタイムのデータ処理、ストレージ、分析が求められるようになりました。バッチとリアルタイムの統合分析は、このニーズに対応します。
バッチとリアルタイムの統合分析は、単一のプラットフォーム上でリアルタイムデータとバッチデータの両方を管理・処理し、リアルタイム処理とバッチ分析のシームレスな連携を可能にします。主な利点は次のとおりです:
-
データ処理効率の向上:リアルタイムデータとバッチデータを単一のプラットフォームに統合することで、データ転送と変換のコストを削減します。
-
分析精度の向上:リアルタイムデータとバッチデータを組み合わせて分析することで、結果の精度が向上します。
-
データ管理の簡素化:統合されたアプローチにより、データ管理と処理が合理化されます。
-
より良い意思決定支援:データを最大限に活用して、ビジネス上の意思決定を支援します。
Alibaba Cloud は、バッチとリアルタイムの両方のシナリオに対応する、簡素化された統合データウェアハウスソリューションを提供します。このソリューションでは、バッチ処理に MaxCompute を、リアルタイム分析に Hologres を使用します。Realtime Compute for Apache Flink のリアルタイム処理能力と組み合わせることで、これらのサービスは Alibaba Cloud の統合データウェアハウスの中核エンジンを形成します。
ソリューションアーキテクチャ
以下の図は、MaxCompute と Hologres を使用して GitHub の公開イベントデータセットに対するバッチとリアルタイムの統合分析を行う完全なパイプラインを示しています。

このアーキテクチャでは、ECS インスタンスがデータソースとして GitHub からリアルタイムおよびバッチのイベントデータを収集・集約します。データはリアルタイムパイプラインとバッチパイプラインに供給され、その後、統合サービスレイヤーとして Hologres に統合されます。
-
リアルタイムパイプライン:Realtime Compute for Apache Flink は、Simple Log Service (SLS) からのデータをリアルタイムで処理し、Hologres に書き込みます。Hologres はリアルタイムのデータ書き込みと更新をサポートし、データは取り込み後すぐにクエリ可能です。これらのネイティブな統合により、最新イベントの抽出やトレンドイベントの分析といったユースケース向けに、高スループット、低レイテンシー、モデル駆動型のリアルタイムデータウェアハウス開発が可能になります。
-
バッチパイプライン:MaxCompute は、大量のバッチデータを処理・アーカイブします。Object Storage Service (OSS) は、生の JSON データに対して、便利で安全、低コストなストレージを提供します。MaxCompute は、外部テーブルを介して OSS 内の半構造化データを直接読み取り・解析し、価値の高いデータを内部ストレージに統合し、DataWorks と連携してバッチデータウェアハウスを構築できます。
-
Hologres はストレージレイヤーで MaxCompute とシームレスに統合されており、MaxCompute 内の大量の履歴データに対するクエリを高速化できます。これにより、履歴データに対する低頻度・パフォーマンス専有型のクエリがサポートされます。また、バッチパイプラインを使用してリアルタイムデータを修正し、リアルタイムパイプラインでのデータ欠落などの問題を解決することもできます。
このソリューションには、以下の利点があります:
-
安定かつ効率的なバッチパイプライン:時間単位のデータ書き込みと更新、大規模なバッチ処理、複雑な計算をサポートし、コンピューティングコストを削減します。
-
成熟したリアルタイムパイプライン:リアルタイムの取り込み、イベント計算、分析をサポートし、数秒以内に応答を返します。
-
統合されたストレージとサービス:Hologres は、一元化されたデータストレージと一貫した外部インターフェイス (OLAP と Key-Value クエリの両方に対応する単一の SQL インターフェイス) を備えた統合サービスレイヤーを提供します。
-
バッチとリアルタイムの統合分析:データの冗長性と移動を削減し、データの修正を可能にします。
このワンストップの開発アプローチにより、秒レベルのデータ応答、エンドツーエンドのステータス可視性、コンポーネントの少ない簡素化されたアーキテクチャ、そして O&M コストの削減が実現します。
ビジネスとデータの理解
開発者は、GitHub 上のオープンソースプロジェクトで作業する際に多くのイベントを作成します。GitHub は、イベントタイプ、開発者、コードリポジトリなど、各イベントの詳細を記録します。GitHub は、リポジトリへのスター付けやコードのコミットなど、公開イベントを利用可能にしています。イベントタイプの完全なリストについては、「Webhook のイベントとペイロード」をご参照ください。
-
GitHub は OpenAPI を通じて公開イベントを提供します。この API は 5 分遅延のリアルタイムデータを提供します。詳細については、「イベント」をご参照ください。
-
GH Archive プロジェクトは、GitHub の公開イベントの時間単位のアーカイブを収集・提供します。これらのアーカイブを使用してオフラインデータを取得します。詳細については、「GH Archive」をご参照ください。
GitHub ビジネスの理解
GitHub の中核ビジネスは、コードとインタラクションの管理です。これには、開発者、リポジトリ、組織という 3 つの主要なエンティティが関わります。
このデータ分析では、Event もエンティティとして保存・記録されます。

公開イベントの生データの理解
以下の例は、ある生イベントの JSON データを示しています:
{
"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"
}
この分析では、15 種類の公開イベントを対象とします。発生しなかったイベントや、記録されなくなったイベントは含まれません。これらのイベントタイプの詳細については、「Github の公開イベントタイプ」をご参照ください。
前提条件
-
Elastic Compute Service (ECS) インスタンスが作成され、Elastic IP アドレス (EIP) が関連付けられていること。このインスタンスは、GitHub API からリアルタイムのイベントデータを抽出するために使用されます。詳細については、「作成ガイド」および「Elastic IP アドレス」をご参照ください。
-
Object Storage Service (OSS) が有効化され、GH Archive からの JSON データファイルを保存するために、ECS インスタンスに ossutil ツールがインストールされていること。詳細については、「OSS の有効化」および「ossutil のインストール」をご参照ください。
-
MaxCompute が有効化され、プロジェクトが作成されていること。詳細については、「MaxCompute プロジェクトの作成」をご参照ください。
-
DataWorks が有効化され、オフラインスケジューリングタスクを作成するためのワークスペースが作成されていること。詳細については、「ワークスペースの作成」をご参照ください。
-
Simple Log Service (SLS) が有効化され、ECS インスタンスからログとしてデータを収集するためのプロジェクトと Logstore が作成されていること。詳細については、「LoongCollector を使用した ECS テキストログの収集と分析」をご参照ください。
-
SLS から Hologres へログデータをリアルタイムで書き込むために、Realtime Compute for Apache Flink インスタンスが有効化されていること。詳細については、「Realtime Compute for Apache Flink の有効化」をご参照ください。
-
Hologres が有効化されていること。詳細については、「Hologres インスタンスの購入」をご参照ください。
オフラインデータウェアハウスの構築 (時間単位の更新)
ECS インスタンスを使用した生データファイルのダウンロードと OSS へのアップロード
Elastic Compute Service (ECS) インスタンスを使用して、GH Archive から JSON データファイルをダウンロードします。
-
wgetコマンドを使用して履歴データをダウンロードします。例えば、wget https://data.gharchive.org/{2012..2022}-{01..12}-{01..31}-{0..23}.json.gzを実行して、2012 年から 2022 年までの時間単位のデータをダウンロードします。 -
毎時生成される新しいデータをダウンロードするには、以下のように時間単位のスケジュールタスクを設定します。
説明-
ECS インスタンスに ossutil がインストールされていることを確認してください。詳細については、「ossutil のインストール」をご参照ください。ossutil のインストールパッケージをダウンロードし、ECS インスタンスにアップロードします。
yum install unzipを実行して unzip ソフトウェアをインストールします。その後、ossutil パッケージを解凍し、実行可能ファイルを/usr/bin/ディレクトリに移動します。 -
ECS インスタンスと同じリージョンに Object Storage Service (OSS) バケットを作成していることを確認してください。カスタムのバケット名を使用できます。この例では、バケット名として
githubeventsを使用します。 -
この例では、ファイルは ECS インスタンスの
/opt/hourlydata/gh_dataディレクトリにダウンロードされます。別のディレクトリを使用することもできます。
-
次のコマンドを実行して、
/opt/hourlydataディレクトリにdownload_code.shという名前のファイルを作成します。cd /opt/hourlydata vim download_code.sh -
iを押して編集モードに入り、次のスクリプトを追加します。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} # データを ./gh_data/ ディレクトリにダウンロードします。別のディレクトリを使用することもできます。 wget ${url} -P ./gh_data/ # gh_data ディレクトリに移動します。 cd gh_data # ダウンロードしたデータを JSON ファイルに解凍します。 gzip -d ${d}.json echo ${d}.json # ルートディレクトリに移動します。 cd /root # ossutil を使用してデータを OSS にアップロードします。 # githubevents OSS バケットに hr=${h} ディレクトリを作成します。 ossutil mkdir oss://githubevents/hr=${h} # /opt/hourlydata/gh_data ディレクトリから OSS にデータをアップロードします。別のディレクトリを使用することもできます。 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! -
Esc キーを押し、
:wqと入力して Enter キーを押し、ファイルを保存して閉じます。 -
次のコマンドを実行して、毎時 10 分に
download_code.shスクリプトを実行します。# 1. 次のコマンドを実行し、I を押して編集モードに入ります。 crontab -e # 2. 次のコマンドを追加します。その後、Esc を押し、:wq と入力して Enter を押して終了します。 10 * * * * cd /opt/hourlydata && sh download_code.sh > download.logスクリプトが実行されると、毎時 10 分に前の 1 時間分の JSON ファイルがダウンロードされます。ファイルは ECS インスタンス上で解凍され、
oss://githubeventsのパスで OSS にアップロードされます。前の 1 時間分のファイルのみを読み取るために、アップロード時に各ファイルのパーティションとして'hr=%Y-%M-%D-%H'という名前のディレクトリが作成されます。これにより、後続のデータ書き込み操作では最新のパーティションからのみファイルが読み取られるようになります。
-
外部テーブルを使用した OSS データの MaxCompute へのインポート
MaxCompute クライアントまたは DataWorks の ODPS SQL ノードで次のコマンドを実行します。詳細については、「クライアント (odpscmd) を使用した MaxCompute への接続」または「ODPS SQL タスクの開発」をご参照ください。
-
OSS に保存されている JSON ファイルを読み取るために、外部テーブル
githubeventsを作成します: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/' ;MaxCompute で OSS データにアクセスするための外部テーブルの作成に関する詳細については、「OSS の非構造化データへのアクセス」をご参照ください。
-
データを保存するためのファクトテーブル
dwd_github_events_odpsを作成します。次のコードはデータ定義言語 (DDL) 文です:CREATE TABLE IF NOT EXISTS dwd_github_events_odps ( id BIGINT COMMENT 'イベント ID' ,actor_id BIGINT COMMENT 'イベント開始者の ID' ,actor_login STRING COMMENT 'イベント開始者のログイン名' ,repo_id BIGINT COMMENT 'リポジトリ ID' ,repo_name STRING COMMENT 'リポジトリのフルネーム (owner/repository_name 形式)' ,org_id BIGINT COMMENT 'リポジトリが属する組織の ID' ,org_login STRING COMMENT 'リポジトリが属する組織の名前' ,`type` STRING COMMENT 'イベントタイプ' ,created_at DATETIME COMMENT 'イベント発生時刻' ,action STRING COMMENT 'イベントアクション' ,iss_or_pr_id BIGINT COMMENT 'issue または pull request の ID' ,number BIGINT COMMENT 'issue または pull request の番号' ,comment_id BIGINT COMMENT 'コメント ID' ,commit_id STRING COMMENT 'コミット ID' ,member_id BIGINT COMMENT 'メンバー ID' ,rev_or_push_or_rel_id BIGINT COMMENT 'review、push、または release の ID' ,ref STRING COMMENT '作成または削除されたリソースの名前' ,ref_type STRING COMMENT '作成または削除されたリソースのタイプ' ,state STRING COMMENT 'issue、pull request、または pull request review のステータス' ,author_association STRING COMMENT 'アクターとリポジトリの関係' ,language STRING COMMENT 'pull request のコードの言語' ,merged BOOLEAN COMMENT 'pull request がマージされたかどうかを示す' ,merged_at DATETIME COMMENT 'コードがマージされた時刻' ,additions BIGINT COMMENT '追加されたコードの行数' ,deletions BIGINT COMMENT '削除されたコードの行数' ,changed_files BIGINT COMMENT 'pull request で変更されたファイルの数' ,push_size BIGINT COMMENT 'コミット数' ,push_distinct_size BIGINT COMMENT '個別コミット数' ,hr STRING COMMENT 'イベントが発生した時間。例えば、イベントが 00:23 に発生した場合、hr の値は 00 です。' ,`month` STRING COMMENT 'イベントが発生した月。例えば、イベントが 2015 年 10 月に発生した場合、month の値は 2015-10 です。' ,`year` STRING COMMENT 'イベントが発生した年。例えば、イベントが 2015 年に発生した場合、year の値は 2015 です。' ) PARTITIONED BY ( ds STRING COMMENT 'イベントが発生した日付 (yyyy-mm-dd 形式)。' ); -
JSON データを解析し、ファクトテーブルに書き込みます。
次のコマンドを実行してパーティションを追加し、JSON データを解析して、
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); -
データのクエリ
次のコマンドを実行して、
dwd_github_events_odpsテーブルからデータをクエリします:SET odps.sql.allow.fullscan=true; SELECT * FROM dwd_github_events_odps where ds = '2023-03-31' limit 10;次のサンプル結果が返されます:
結果には次のフィールドが含まれます:
-
id:イベント ID -
actor_id/actor_login:ユーザー ID とユーザー名 -
repo_id/repo_name:リポジトリ ID とリポジトリ名 -
org_id/org_login:組織 ID と組織名 (一部のイベントでは空) -
type:イベントタイプ (例:CreateEvent、PushEvent、DeleteEvent、PullRequestReviewEvent) -
created_at:作成時刻 -
action:アクションタイプ
-
リアルタイムデータウェアハウスの構築
ECS を使用したリアルタイムデータの取得
Elastic Compute Service (ECS) インスタンスを使用して、GitHub API からリアルタイムのイベントデータを抽出します。次のサンプルスクリプトは、GitHub API からリアルタイムデータを収集する方法を示しています。
-
スクリプトが実行されるたびに、1 分間実行されます。この期間中に API が提供するリアルタイムのイベントデータを収集し、各イベントを JSON 形式で保存します。
-
このスクリプトは、すべてのリアルタイムイベントデータが収集されることを保証するものではありません。
-
GitHub API から継続的にデータを収集するには、Accept ヘッダーと Authorization ヘッダーを提供する必要があります。Accept の値は固定です。Authorization には、GitHub から取得した個人アクセストークンを入力します。個人アクセストークンの作成方法の詳細については、「こちらのドキュメント」をご参照ください。
-
次のコマンドを実行して、
/opt/realtimeディレクトリにdownload_realtime_data.pyという名前のファイルを作成します。cd /opt/realtime vim download_realtime_data.py -
iを押して編集モードに入り、次のサンプルコンテンツをファイルに追加します。#!python import requests import json import sys import time # 次のページの 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 # API から 1 ページ分のデータを収集 def download(link, fname): # GitHub API の Accept ヘッダーと Authorization ヘッダーを定義 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 # 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 # 現在時刻を定義 def get_current_ms(): return round(time.time()*1000) # スクリプトの実行時間を 1 分と定義 def main(fname): current_ms = get_current_ms() while get_current_ms() - current_ms < 60*1000: download_all_data(fname) time.sleep(0.1) # スクリプトを実行 if __name__ == '__main__': if len(sys.argv) < 2: print('usage: python {} <log_file>'.format(sys.argv[0])) exit(0) main(sys.argv[1]) -
Esc キーを押し、
:wqと入力して Enter キーを押し、ファイルを保存して閉じます。 -
download_realtime_data.pyを実行し、各実行で収集したデータを個別に保存するためにrun_py.shファイルを作成します。内容は次のとおりです。python /opt/realtime/download_realtime_data.py /opt/realtime/gh_realtime_data/$(date '+%Y-%m-%d-%H:%M:%S').json -
履歴データを削除するために
delete_log.shファイルを作成します。内容は次のとおりです。d=$(TZ=UTC date --date='2 day ago' '+%Y-%m-%d') rm -f /opt/realtime/gh_realtime_data/*${d}*.json -
次のコマンドを実行して、毎分 GitHub データを収集し、毎日履歴データを削除します。
#1. 次のコマンドを実行し、I を押して編集モードに入ります。 crontab -e #2. 次のコマンドを追加します。その後、Esc を押し、:wq と入力して Enter を押して終了します。 * * * * * bash /opt/realtime/run_py.sh 1 1 * * * bash /opt/realtime/delete_log.sh
SLS を使用した ECS データの収集
Simple Log Service (SLS) は、ECS インスタンスからのリアルタイムイベントデータをログとして収集します。
SLS は Logtail を使用して ECS インスタンスからログを収集することをサポートしています。データは JSON 形式であるため、Logtail の JSON モードを使用して、ECS インスタンスから増分 JSON ログを迅速に収集できます。詳細については、「JSON モードでのログ収集」をご参照ください。このトピックでは、SLS は生データのトップレベルのキーと値のペアを解析するように構成されています。
この例では、Logtail 構成のログパスパラメーターは /opt/realtime/gh_realtime_data/**/*.json に設定されています。
構成が完了すると、SLS は ECS インスタンスから増分イベントデータを継続的に収集します。収集されたログデータは、SLS コンソールの [生ログ] タブで表示できます。各ログエントリには、actor、created_at、id、org、payload、public、repo、type などの解析されたトップレベルのフィールドが含まれています。
Flink を使用した SLS データの Hologres へのリアルタイム書き込み
Flink は、SLS によって収集されたログデータをリアルタイムで Hologres に書き込みます。Flink で SLS ソーステーブルと Hologres 結果テーブルを使用することで、SLS から Hologres にデータをストリーミングできます。詳細については、「Simple Log Service からのデータインポート」をご参照ください。
-
Hologres 内部テーブルの作成
内部テーブルには、生の JSON データからの一部のキーと値のペアのみが保持されます。イベント
idと日付dsがプライマリキーとして設定されます。イベントidは分散キーとして設定されます。日付dsはパーティションキーとして設定されます。イベント時間created_atは event_time_column として設定されます。必要に応じて他のフィールドにインデックスを作成できます。インデックスの詳細については、「CREATE TABLE」をご参照ください。この例では、テーブルを作成するために次のデータ定義言語 (DDL) 文が使用されます。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 'イベント ID'; COMMENT ON COLUMN public.gh_realtime_data.actor_id IS 'イベント開始者の ID'; COMMENT ON COLUMN public.gh_realtime_data.actor_login IS 'イベント開始者のログイン名'; COMMENT ON COLUMN public.gh_realtime_data.repo_id IS 'repo ID'; COMMENT ON COLUMN public.gh_realtime_data.repo_name IS 'repo 名'; COMMENT ON COLUMN public.gh_realtime_data.org_id IS 'repo が属する組織の ID'; COMMENT ON COLUMN public.gh_realtime_data.org_login IS 'repo が属する組織の名前'; COMMENT ON COLUMN public.gh_realtime_data.type IS 'イベントタイプ'; COMMENT ON COLUMN public.gh_realtime_data.created_at IS 'イベント発生時刻'; COMMENT ON COLUMN public.gh_realtime_data.action IS 'イベントアクション'; 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 番号'; COMMENT ON COLUMN public.gh_realtime_data.comment_id IS 'コメント ID'; COMMENT ON COLUMN public.gh_realtime_data.commit_id IS 'コミット ID'; COMMENT ON COLUMN public.gh_realtime_data.member_id IS 'メンバー 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 '作成または削除されたリソースの名前'; COMMENT ON COLUMN public.gh_realtime_data.ref_type IS '作成または削除されたリソースのタイプ'; COMMENT ON COLUMN public.gh_realtime_data.state IS 'issue/pull_request/pull_request_review のステータス'; COMMENT ON COLUMN public.gh_realtime_data.author_association IS 'アクターと repo の関係'; COMMENT ON COLUMN public.gh_realtime_data.language IS 'プログラミング言語'; COMMENT ON COLUMN public.gh_realtime_data.merged IS 'マージが受け入れられたかどうかを指定します'; COMMENT ON COLUMN public.gh_realtime_data.merged_at IS 'コードがマージされた時刻'; COMMENT ON COLUMN public.gh_realtime_data.additions IS '追加されたコードの行数'; COMMENT ON COLUMN public.gh_realtime_data.deletions IS '削除されたコードの行数'; COMMENT ON COLUMN public.gh_realtime_data.changed_files IS 'pull request で変更されたファイルの数'; COMMENT ON COLUMN public.gh_realtime_data.push_size IS 'プッシュ数'; COMMENT ON COLUMN public.gh_realtime_data.push_distinct_size IS '個別プッシュ数'; COMMENT ON COLUMN public.gh_realtime_data.hr IS 'イベントが発生した時間。例えば、時刻が 00:23 の場合、hr=00。'; COMMENT ON COLUMN public.gh_realtime_data.month IS 'イベントが発生した月。例えば、日付が 2015 年 10 月の場合、month=2015-10。'; COMMENT ON COLUMN public.gh_realtime_data.year IS 'イベントが発生した年。例えば、年が 2015 年の場合、year=2015。'; COMMENT ON COLUMN public.gh_realtime_data.ds IS 'イベントが発生した日。ds=yyyy-mm-dd。'; COMMIT; -
Flink を使用したリアルタイムでのデータ書き込み
Flink を使用して SLS データを解析し、リアルタイムで Hologres に書き込みます。次の Flink 文はデータをフィルタリングします:イベント ID またはイベント時間 (
created_at) が null のダーティデータは破棄され、最近のイベントデータのみが保持されます。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>',--SLS のプライベートエンドポイント 'accessid' = '<accesskey id>',--アカウントの AccessKey ID 'accesskey' = '<accesskey secret>',--アカウントの AccessKey Secret 'project' = '<project name>',--SLS プロジェクトの名前 'logstore' = '<logstore name>'--SLS Logstore の名前 'starttime' = '2023-04-06 00:00:00',--SLS データ収集の開始時刻 ); 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>', --Hologres データベースの名前 'tablename' = '<hologres tablename>', --データを受け取る Hologres テーブルの名前 'username' = '<accesskey id>', --現在の Alibaba Cloud アカウントの AccessKey ID 'password' = '<accesskey secret>', --現在の Alibaba Cloud アカウントの AccessKey Secret 'endpoint' = '<endpoint>', --現在の Hologres インスタンスの VPC エンドポイント 'jdbcretrycount' = '1', --接続失敗時のリトライ回数 'partitionrouter' = 'true', --パーティションテーブルにデータを書き込むかどうかを指定 'createparttable' = 'true', --パーティションを自動的に作成するかどうかを指定 'mutatetype' = 'insertorignore' --データ書き込みモード ); 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);パラメーターの詳細については、「Simple Log Service (SLS)」および「Hologres」をご参照ください。
説明GitHub からの生のイベントデータは UTC タイムゾーンを使用し、タイムゾーン属性を持ちません。Hologres のデフォルトのタイムゾーンは UTC+8 です。したがって、Flink から Hologres にリアルタイムでデータを書き込む際にタイムゾーンを調整する必要があります。Flink SQL でソーステーブルデータに UTC タイムゾーン属性を割り当てます。手順は次のとおりです:
ステップ 1:ジョブ編集ページに移動
-
Realtime Compute for Apache Flink コンソールにログインします
-
対象のワークスペースに移動します
-
対象の Flink SQL または JAR ジョブを見つけて、[編集] をクリックします
ステップ 2:[デプロイメント詳細] タブを開く
⚠️ 重要:「Flink 構成」は独立したセクションではなくなり、「デプロイメント詳細」に統合されました。
-
ジョブ編集ページの上部で、[デプロイメント詳細] タブに切り替えます。
-
[パラメーター構成] セクションまでスクロールダウンします。
ステップ 3:カスタム構成の追加
-
[パラメーター構成] の右側にある [編集] ボタンをクリックします。
-
表示されるダイアログボックスで、[その他の構成] テキストボックスを見つけます。
-
テキストボックスに、Flink パラメーター
table.local-time-zone:Asia/Shanghaiをキーと値のペアとして追加し、Flink システムタイムゾーンをAsia/Shanghaiに設定します。
-
-
データのクエリ
Flink を介して Hologres に書き込まれた SLS データをクエリし、必要に応じてデータ開発を行います。
SELECT * FROM public.gh_realtime_data limit 10;クエリ結果は次のフィールドを返します:
-
id -
actor_id -
actor_login -
repo_id -
repo_name -
org_id -
org_login -
type(イベントタイプ、例:PullRequestReviewEvent、CreateEvent、PushEvent、PullRequestEvent) -
created_at -
action -
iss_or_pr_id
-
オフラインデータを使用したリアルタイムデータの修正
このシナリオでは、リアルタイムデータが欠落する可能性があります。オフラインデータを使用してリアルタイムデータを修正できます。次の手順は、前日のリアルタイムデータを修正する方法を示しています。必要に応じて修正期間を調整してください。
-
Hologres で外部テーブルを作成して、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');パラメーターの詳細については、「IMPORT FOREIGN SCHEMA」をご参照ください。
-
一時テーブルを作成して、前日のリアルタイムデータをオフラインデータで修正します。
説明Hologres V2.1.17 以降では、サーバーレスコンピューティングがサポートされています。大規模なオフラインデータインポート、大規模な ETL ジョブ、外部テーブルに対する大量のクエリなどのシナリオでは、サーバーレスコンピューティングを使用してこれらのタスクを実行できます。この機能は、インスタンスリソースの代わりに、追加のサーバーレスリソースを使用するため、インスタンスの安定性が向上し、メモリ不足 (OOM) エラーの可能性が減少します。インスタンスに追加のコンピューティングリソースを予約する必要はなく、実行したタスクに対してのみ課金されます。サーバーレスコンピューティングの詳細については、「サーバーレスコンピューティング」をご参照ください。サーバーレスコンピューティングの使用方法については、「サーバーレスコンピューティングの使用」をご参照ください。
-- 潜在的な一時テーブルをクリーンアップ DROP TABLE IF EXISTS gh_realtime_data_tmp; -- 一時テーブルを作成 SET hg_experimental_enable_create_table_like_properties = ON; CALL HG_CREATE_TABLE_LIKE ('gh_realtime_data_tmp', 'select * from gh_realtime_data'); -- (オプション) サーバーレスコンピューティングを使用して、大規模なオフラインデータインポートと ETL ジョブを実行 SET hg_computing_resource = 'serverless'; -- 一時テーブルにデータを挿入し、統計を更新 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; -- 必須でない SQL 文がサーバーレスリソースを使用しないように構成をリセット RESET hg_computing_resource; -- アトミックテーブルを既存の一時子テーブルに置き換え 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;
データ分析
収集したデータに対して、幅広い分析を実行できます。ビジネスで必要な時間範囲に基づいて、リアルタイム分析、オフライン分析、およびリアルタイムとオフラインの統合分析をサポートするように、データウェアハウスをレイヤー化して設計します。
以下の例では、リアルタイムデータを分析します。特定のコードリポジトリや開発者に関するデータを分析することもできます。
-
本日の公開イベントの総数をクエリします。
SELECT count(*) FROM gh_realtime_data WHERE created_at >= date_trunc('day', now());以下はサンプル結果です:
count ------ 1006 -
過去 1 日で最もアクティブなプロジェクト (最もイベントが多い) をクエリします。
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;以下はサンプル結果です:
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 -
過去 1 日で最もアクティブな開発者 (最もイベントが多い) をクエリします。
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;以下はサンプル結果です:
actor_login events ------------------+------ direwolf-github 13 arm-on 10 sergii-nosachenko 9 Christoffel-T 9 yangwang201911 7 -
過去 1 時間で最も人気のあるプログラミング言語のランキングをクエリします。
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;以下はサンプル結果です:
language total -----------+---- JavaScript 25 C++ 15 Python 14 TypeScript 13 Java 8 PHP 8 -
過去 1 日に受け取ったスターの数によるプロジェクトのランキングをクエリします。
説明この例では、ユーザーがプロジェクトのスターを外すケースは考慮されていません。
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;以下はサンプル結果です:
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 -
本日のデイリーアクティブユーザーとプロジェクトをクエリします。
SELECT uniq (actor_id) actor_num, uniq (repo_id) repo_num FROM gh_realtime_data WHERE created_at > date_trunc('day', now());以下はサンプル結果です:
actor_num repo_num ---------+-------- 743 816