Bangun solusi analitik batch dan real-time terpadu menggunakan data event GitHub. MaxCompute berperan sebagai gudang data batch, sedangkan Realtime Compute for Apache Flink dan Hologres membentuk gudang data real-time. Hologres dan MaxCompute kemudian menyediakan lapisan terpadu untuk analisis data real-time maupun batch.
Latar Belakang
Seiring bisnis semakin digital, permintaan akan data yang lebih mutakhir terus meningkat. Di luar pemrosesan batch tradisional untuk data skala besar, banyak bisnis kini memerlukan pemrosesan, penyimpanan, dan analitik data secara real-time. Analitik batch dan real-time terpadu menjawab kebutuhan ini.
Analitik batch dan real-time terpadu mengelola serta memproses data real-time dan batch pada satu platform yang sama, sehingga menciptakan koneksi mulus antara pemrosesan real-time dan analitik batch. Manfaat utamanya meliputi:
-
Peningkatan efisiensi pemrosesan data: Integrasi data real-time dan batch pada satu platform mengurangi biaya transfer dan konversi data.
-
Peningkatan akurasi analitik: Menggabungkan data real-time dan batch untuk analisis meningkatkan presisi hasil Anda.
-
Penyederhanaan manajemen data: Pendekatan terpadu menyederhanakan manajemen dan pemrosesan data.
-
Dukungan pengambilan keputusan yang lebih baik: Manfaatkan sepenuhnya data Anda untuk mendukung keputusan bisnis.
Alibaba Cloud menawarkan solusi gudang data terpadu yang disederhanakan untuk skenario batch maupun real-time. Solusi ini menggunakan MaxCompute untuk pemrosesan batch dan Hologres untuk analitik real-time. Dipadukan dengan kemampuan pemrosesan real-time dari Realtime Compute for Apache Flink, layanan-layanan ini membentuk mesin inti gudang data terpadu Alibaba Cloud.
Arsitektur solusi
Diagram berikut menunjukkan pipeline lengkap untuk analitik batch dan real-time terpadu pada dataset event publik GitHub menggunakan MaxCompute dan Hologres.

Dalam arsitektur ini, sebuah instans ECS mengumpulkan dan mengagregasi data event real-time dan batch dari GitHub sebagai sumber data. Data tersebut dialirkan ke pipeline real-time dan pipeline batch, lalu dikonsolidasikan di Hologres sebagai lapisan layanan terpadu.
-
Pipeline real-time: Realtime Compute for Apache Flink memproses data dari Simple Log Service (SLS) secara real-time dan menuliskannya ke Hologres. Hologres mendukung penulisan dan pembaruan data real-time, dengan data dapat langsung diquery setelah diingest. Integrasi native mereka memungkinkan pengembangan gudang data real-time berbasis model dengan throughput tinggi dan latensi rendah untuk kasus penggunaan seperti mengekstraksi event terbaru dan menganalisis event yang sedang tren.
-
Pipeline batch: MaxCompute memproses dan mengarsipkan data batch dalam jumlah besar. Object Storage Service (OSS) menyediakan penyimpanan mentah JSON yang nyaman, aman, dan berbiaya rendah. MaxCompute dapat langsung membaca dan mengurai data semi-terstruktur di OSS melalui tabel eksternal, mengintegrasikan data bernilai tinggi ke penyimpanan internalnya, serta bekerja sama dengan DataWorks untuk membangun gudang data batch.
-
Hologres terintegrasi secara mulus dengan MaxCompute di lapisan penyimpanan, memungkinkan Anda mempercepat query pada volume besar data historis di MaxCompute. Hal ini mendukung query frekuensi rendah namun berkinerja tinggi pada data historis. Anda juga dapat menggunakan pipeline batch untuk mengoreksi data real-time dan menyelesaikan isu seperti kelalaian data pada pipeline real-time.
Solusi ini menawarkan keunggulan berikut:
-
Pipeline batch yang stabil dan efisien: Mendukung penulisan dan pembaruan data per jam, pemrosesan batch skala besar, perhitungan kompleks, serta pengurangan biaya komputasi.
-
Pipeline real-time yang matang: Mendukung ingest real-time, komputasi event, dan analitik, memberikan respons dalam hitungan detik.
-
Penyimpanan dan layanan terpadu: Hologres menyediakan lapisan layanan terpadu dengan penyimpanan data terpusat dan antarmuka eksternal yang konsisten (satu antarmuka SQL untuk query OLAP dan Key-Value).
-
Analitik batch dan real-time terpadu: Mengurangi redundansi dan perpindahan data, serta memungkinkan koreksi data.
Pendekatan pengembangan one-stop ini mencapai respons data tingkat detik, visibilitas status end-to-end, arsitektur yang disederhanakan dengan komponen lebih sedikit, serta pengurangan biaya O&M.
Memahami bisnis dan data
Developer membuat banyak event saat mengerjakan proyek open source di GitHub. GitHub mencatat detail setiap event, termasuk jenis event, developer, dan repositori kode. GitHub menyediakan event publik, seperti memberi bintang pada repositori atau melakukan commit kode. Untuk daftar lengkap jenis event, lihat Webhook events and payloads.
-
GitHub menyediakan event publik melalui OpenAPI. API ini menawarkan data real-time dengan keterlambatan lima menit. Untuk informasi selengkapnya, lihat Events.
-
Proyek GH Archive mengumpulkan dan menyediakan arsip per jam event publik GitHub. Gunakan arsip ini untuk mendapatkan data offline. Untuk informasi selengkapnya, lihat GH Archive.
Memahami bisnis GitHub
Bisnis inti GitHub adalah mengelola kode dan interaksi. Ini melibatkan tiga entitas utama: Developer, Repository, dan Organization.
Untuk analisis data ini, sebuah Event juga disimpan dan dicatat sebagai entitas.

Memahami data event publik mentah
Contoh berikut menunjukkan data JSON untuk sebuah event mentah:
{
"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"
}
Analisis ini mencakup 15 jenis event publik. Tidak termasuk event yang tidak pernah terjadi atau tidak lagi direkam. Untuk detail tentang jenis-jenis event ini, lihat Github public event types.
Prasyarat
-
Sebuah instans Elastic Compute Service (ECS) telah dibuat dan Elastic IP Address (EIP) telah dikaitkan dengannya. Instans ini digunakan untuk mengekstraksi data event real-time dari GitHub API. Untuk informasi selengkapnya, lihat Panduan pembuatan dan Elastic IP Address.
-
Object Storage Service (OSS) telah diaktifkan, dan tool ossutil telah diinstal pada instans ECS untuk menyimpan file data JSON dari GH Archive. Untuk informasi selengkapnya, lihat Aktifkan OSS dan Instal ossutil.
-
MaxCompute telah diaktifkan dan sebuah proyek telah dibuat. Untuk informasi selengkapnya, lihat Buat proyek MaxCompute.
-
DataWorks telah diaktifkan dan sebuah ruang kerja telah dibuat untuk membuat tugas penjadwalan offline. Untuk informasi selengkapnya, lihat Buat ruang kerja.
-
Simple Log Service (SLS) telah diaktifkan, dan sebuah proyek serta Logstore telah dibuat untuk mengumpulkan data dari instans ECS sebagai log. Untuk informasi selengkapnya, lihat Kumpulkan dan analisis log teks ECS menggunakan LoongCollector.
-
Sebuah instans Realtime Compute for Apache Flink telah diaktifkan untuk menuliskan data log dari SLS ke Hologres secara real-time. Untuk informasi selengkapnya, lihat Aktifkan Realtime Compute for Apache Flink.
-
Hologres telah diaktifkan. Untuk informasi selengkapnya, lihat Beli instans Hologres.
Bangun gudang data offline (pembaruan per jam)
Unduh file data mentah menggunakan instans ECS dan unggah ke OSS
Gunakan instans Elastic Compute Service (ECS) untuk mengunduh file data JSON dari GH Archive.
-
Unduh data historis menggunakan perintah
wget. Misalnya, jalankanwget https://data.gharchive.org/{2012..2022}-{01..12}-{01..31}-{0..23}.json.gzuntuk mengunduh data per jam dari tahun 2012 hingga 2022. -
Untuk mengunduh data baru yang dihasilkan setiap jam, atur tugas terjadwal per jam sebagai berikut.
Catatan-
Pastikan ossutil telah diinstal pada instans ECS. Untuk informasi selengkapnya, lihat Instal ossutil. Unduh paket instalasi ossutil dan unggah ke instans ECS. Jalankan
yum install unzipuntuk menginstal software unzip. Kemudian, dekompresi paket ossutil dan pindahkan file yang dapat dieksekusi ke direktori/usr/bin/. -
Pastikan Anda telah membuat bucket Object Storage Service (OSS) di wilayah yang sama dengan instans ECS Anda. Anda dapat menggunakan nama bucket kustom. Contoh ini menggunakan nama bucket
githubevents. -
Dalam contoh ini, file diunduh ke direktori
/opt/hourlydata/gh_datapada instans ECS. Anda dapat menggunakan direktori berbeda.
-
Jalankan perintah berikut untuk membuat file bernama
download_code.shdi direktori/opt/hourlydata.cd /opt/hourlydata vim download_code.sh -
Tekan
iuntuk masuk ke mode edit dan tambahkan skrip berikut.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} # Unduh data ke direktori ./gh_data/. Anda dapat menggunakan direktori berbeda. wget ${url} -P ./gh_data/ # Pindah ke direktori gh_data. cd gh_data # Dekompresi data yang diunduh menjadi file JSON. gzip -d ${d}.json echo ${d}.json # Pindah ke direktori root. cd /root # Gunakan ossutil untuk mengunggah data ke OSS. # Buat direktori hr=${h} di bucket OSS githubevents. ossutil mkdir oss://githubevents/hr=${h} # Unggah data dari direktori /opt/hourlydata/gh_data ke OSS. Anda dapat menggunakan direktori berbeda. 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! -
Tekan tombol Esc, ketik
:wq, lalu tekan Enter untuk menyimpan dan menutup file. -
Jalankan perintah berikut untuk mengeksekusi skrip
download_code.shpada menit ke-10 setiap jam.# 1. Jalankan perintah berikut dan tekan I untuk masuk ke mode edit. crontab -e # 2. Tambahkan perintah berikut. Lalu, tekan Esc, ketik :wq, dan tekan Enter untuk keluar. 10 * * * * cd /opt/hourlydata && sh download_code.sh > download.logSetelah skrip dijalankan, file JSON dari jam sebelumnya diunduh pada menit ke-10 setiap jam. File tersebut kemudian didekompresi pada instans ECS dan diunggah ke OSS di path
oss://githubevents. Untuk hanya membaca file dari jam sebelumnya, direktori bernama'hr=%Y-%M-%D-%H'dibuat sebagai partisi untuk setiap file selama pengunggahan. Hal ini memastikan bahwa operasi penulisan data selanjutnya hanya membaca file dari partisi terbaru.
-
Impor data OSS ke MaxCompute menggunakan tabel eksternal
Jalankan perintah berikut di klien MaxCompute atau node ODPS SQL di DataWorks. Untuk informasi selengkapnya, lihat Hubungkan ke MaxCompute menggunakan klien (odpscmd) atau Kembangkan tugas ODPS SQL.
-
Buat tabel eksternal
githubeventsuntuk membaca file JSON yang disimpan di 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/' ;Untuk informasi selengkapnya tentang membuat tabel eksternal untuk mengakses data OSS di MaxCompute, lihat Akses data tidak terstruktur di OSS.
-
Buat tabel fakta
dwd_github_events_odpsuntuk menyimpan data. Kode berikut menunjukkan pernyataan Data Definition Language (DDL):CREATE TABLE IF NOT EXISTS dwd_github_events_odps ( id BIGINT COMMENT 'Event ID' ,actor_id BIGINT COMMENT 'ID dari inisiator event' ,actor_login STRING COMMENT 'Nama login dari inisiator event' ,repo_id BIGINT COMMENT 'ID repositori' ,repo_name STRING COMMENT 'Nama lengkap repositori dalam format owner/repository_name' ,org_id BIGINT COMMENT 'ID organisasi tempat repositori berada' ,org_login STRING COMMENT 'Nama organisasi tempat repositori berada' ,`type` STRING COMMENT 'Jenis event' ,created_at DATETIME COMMENT 'Waktu terjadinya event' ,action STRING COMMENT 'Aksi event' ,iss_or_pr_id BIGINT COMMENT 'ID issue atau pull request' ,number BIGINT COMMENT 'Nomor issue atau pull request' ,comment_id BIGINT COMMENT 'ID komentar' ,commit_id STRING COMMENT 'ID commit' ,member_id BIGINT COMMENT 'ID anggota' ,rev_or_push_or_rel_id BIGINT COMMENT 'ID review, push, atau release' ,ref STRING COMMENT 'Nama resource yang dibuat atau dihapus' ,ref_type STRING COMMENT 'Jenis resource yang dibuat atau dihapus' ,state STRING COMMENT 'Status issue, pull request, atau review pull request' ,author_association STRING COMMENT 'Hubungan antara actor dan repositori' ,language STRING COMMENT 'Bahasa kode dalam pull request' ,merged BOOLEAN COMMENT 'Menunjukkan apakah pull request telah digabung' ,merged_at DATETIME COMMENT 'Waktu kode digabung' ,additions BIGINT COMMENT 'Jumlah baris kode yang ditambahkan' ,deletions BIGINT COMMENT 'Jumlah baris kode yang dihapus' ,changed_files BIGINT COMMENT 'Jumlah file yang diubah dalam pull request' ,push_size BIGINT COMMENT 'Jumlah commit' ,push_distinct_size BIGINT COMMENT 'Jumlah commit berbeda' ,hr STRING COMMENT 'Jam terjadinya event. Misalnya, jika event terjadi pukul 00:23, nilai hr adalah 00.' ,`month` STRING COMMENT 'Bulan terjadinya event. Misalnya, jika event terjadi pada Oktober 2015, nilai month adalah 2015-10.' ,`year` STRING COMMENT 'Tahun terjadinya event. Misalnya, jika event terjadi pada 2015, nilai year adalah 2015.' ) PARTITIONED BY ( ds STRING COMMENT 'Tanggal terjadinya event, dalam format yyyy-mm-dd.' ); -
Urai data JSON dan tulis ke tabel fakta.
Jalankan perintah berikut untuk menambahkan partisi, mengurai data JSON, dan menuliskan data ke tabel
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); -
Kueri data.
Jalankan perintah berikut untuk mengquery data dari tabel
dwd_github_events_odps:SET odps.sql.allow.fullscan=true; SELECT * FROM dwd_github_events_odps where ds = '2023-03-31' limit 10;Hasil sampel berikut dikembalikan:
Hasil berisi bidang-bidang berikut:
-
id: ID Event -
actor_id/actor_login: ID pengguna dan username -
repo_id/repo_name: ID repositori dan nama repositori -
org_id/org_login: ID organisasi dan nama organisasi (kosong untuk beberapa event) -
type: Jenis event, seperti CreateEvent, PushEvent, DeleteEvent, PullRequestReviewEvent -
created_at: Waktu pembuatan -
action: Jenis aksi
-
Bangun gudang data real-time
Dapatkan data real-time menggunakan ECS
Sebuah instans Elastic Compute Service (ECS) digunakan untuk mengekstraksi data event real-time dari GitHub API. Skrip sampel berikut menunjukkan cara mengumpulkan data real-time dari GitHub API.
-
Setiap kali skrip dijalankan, ia berjalan selama 1 menit. Skrip ini mengumpulkan data event real-time yang disediakan oleh API selama periode tersebut dan menyimpan setiap event dalam format JSON.
-
Skrip ini tidak menjamin bahwa semua data event real-time dikumpulkan.
-
Untuk terus-menerus mengumpulkan data dari GitHub API, berikan header Accept dan Authorization. Nilai Accept bersifat tetap. Untuk Authorization, masukkan token akses pribadi yang Anda peroleh dari GitHub. Untuk informasi selengkapnya tentang cara membuat token akses pribadi, lihat dokumentasi ini.
-
Jalankan perintah berikut untuk membuat file bernama
download_realtime_data.pydi direktori/opt/realtime.cd /opt/realtime vim download_realtime_data.py -
Tekan
iuntuk masuk ke mode edit dan tambahkan konten sampel berikut ke file.#!python import requests import json import sys import time # Dapatkan URL API 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 # Kumpulkan satu halaman data dari API def download(link, fname): # Definisikan header Accept dan Authorization untuk 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 # Kumpulkan beberapa halaman data dari 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 # Definisikan waktu saat ini def get_current_ms(): return round(time.time()*1000) # Definisikan durasi eksekusi skrip sebagai 1 menit def main(fname): current_ms = get_current_ms() while get_current_ms() - current_ms < 60*1000: download_all_data(fname) time.sleep(0.1) # Jalankan skrip if __name__ == '__main__': if len(sys.argv) < 2: print('usage: python {} <log_file>'.format(sys.argv[0])) exit(0) main(sys.argv[1]) -
Tekan tombol Esc, ketik
:wq, lalu tekan Enter untuk menyimpan dan menutup file. -
Buat file
run_py.shuntuk menjalankandownload_realtime_data.pydan menyimpan data yang dikumpulkan dari setiap eksekusi secara terpisah. Isinya sebagai berikut.python /opt/realtime/download_realtime_data.py /opt/realtime/gh_realtime_data/$(date '+%Y-%m-%d-%H:%M:%S').json -
Buat file
delete_log.shuntuk menghapus data historis. Isinya sebagai berikut.d=$(TZ=UTC date --date='2 day ago' '+%Y-%m-%d') rm -f /opt/realtime/gh_realtime_data/*${d}*.json -
Jalankan perintah berikut untuk mengumpulkan data GitHub setiap menit dan menghapus data historis setiap hari.
#1. Jalankan perintah berikut dan tekan I untuk masuk ke mode edit. crontab -e #2. Tambahkan perintah berikut. Lalu, tekan Esc, ketik :wq, dan tekan Enter untuk keluar. * * * * * bash /opt/realtime/run_py.sh 1 1 * * * bash /opt/realtime/delete_log.sh
Kumpulkan data ECS menggunakan SLS
Simple Log Service (SLS) mengumpulkan data event real-time dari instans ECS sebagai log.
SLS mendukung pengumpulan log dari instans ECS menggunakan Logtail. Karena datanya dalam format JSON, Anda dapat menggunakan mode JSON Logtail untuk dengan cepat mengumpulkan log JSON inkremental dari instans ECS. Untuk informasi selengkapnya, lihat Kumpulkan log dalam mode JSON. Dalam topik ini, SLS dikonfigurasi untuk mengurai pasangan kunci-nilai tingkat atas dari data mentah.
Dalam contoh ini, parameter path log untuk konfigurasi Logtail diatur ke /opt/realtime/gh_realtime_data/**/*.json.
Setelah konfigurasi selesai, SLS terus-menerus mengumpulkan data event inkremental dari instans ECS. Anda dapat melihat data log yang dikumpulkan di tab Raw Logs konsol SLS. Setiap entri log berisi bidang tingkat atas yang telah diurai, seperti actor, created_at, id, org, payload, public, repo, dan type.
Tuliskan data SLS ke Hologres secara real-time menggunakan Flink
Flink menuliskan data log yang dikumpulkan oleh SLS ke Hologres secara real-time. Dengan menggunakan tabel sumber SLS dan tabel hasil Hologres di Flink, Anda dapat mengalirkan data dari SLS ke Hologres. Untuk informasi selengkapnya, lihat Impor data dari Simple Log Service.
-
Buat tabel internal Hologres.
Tabel internal hanya menyimpan beberapa pasangan kunci-nilai dari data JSON mentah.
idevent dan tanggaldsditetapkan sebagai primary key.idevent ditetapkan sebagai distribution key. Tanggaldsditetapkan sebagai partition key. Waktu eventcreated_atditetapkan sebagai event_time_column. Anda dapat membuat indeks untuk bidang lain sesuai kebutuhan. Untuk informasi selengkapnya tentang indeks, lihat CREATE TABLE. Pernyataan Data Definition Language (DDL) berikut digunakan untuk membuat tabel dalam contoh ini.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 dari inisiator event'; COMMENT ON COLUMN public.gh_realtime_data.actor_login IS 'Nama login dari inisiator event'; COMMENT ON COLUMN public.gh_realtime_data.repo_id IS 'ID repositori'; COMMENT ON COLUMN public.gh_realtime_data.repo_name IS 'Nama repositori'; COMMENT ON COLUMN public.gh_realtime_data.org_id IS 'ID organisasi tempat repositori berada'; COMMENT ON COLUMN public.gh_realtime_data.org_login IS 'Nama organisasi tempat repositori berada'; COMMENT ON COLUMN public.gh_realtime_data.type IS 'Jenis event'; COMMENT ON COLUMN public.gh_realtime_data.created_at IS 'Waktu terjadinya event'; COMMENT ON COLUMN public.gh_realtime_data.action IS 'Aksi event'; COMMENT ON COLUMN public.gh_realtime_data.iss_or_pr_id IS 'ID issue/pull_request'; COMMENT ON COLUMN public.gh_realtime_data.number IS 'Nomor issue/pull_request'; COMMENT ON COLUMN public.gh_realtime_data.comment_id IS 'ID komentar'; COMMENT ON COLUMN public.gh_realtime_data.commit_id IS 'ID commit'; COMMENT ON COLUMN public.gh_realtime_data.member_id IS 'ID anggota'; COMMENT ON COLUMN public.gh_realtime_data.rev_or_push_or_rel_id IS 'ID review/push/release'; COMMENT ON COLUMN public.gh_realtime_data.ref IS 'Nama resource yang dibuat atau dihapus'; COMMENT ON COLUMN public.gh_realtime_data.ref_type IS 'Jenis resource yang dibuat atau dihapus'; COMMENT ON COLUMN public.gh_realtime_data.state IS 'Status issue/pull_request/pull_request_review'; COMMENT ON COLUMN public.gh_realtime_data.author_association IS 'Hubungan antara actor dan repositori'; COMMENT ON COLUMN public.gh_realtime_data.language IS 'Bahasa pemrograman'; COMMENT ON COLUMN public.gh_realtime_data.merged IS 'Menentukan apakah merge diterima'; COMMENT ON COLUMN public.gh_realtime_data.merged_at IS 'Waktu kode digabung'; COMMENT ON COLUMN public.gh_realtime_data.additions IS 'Jumlah baris kode yang ditambahkan'; COMMENT ON COLUMN public.gh_realtime_data.deletions IS 'Jumlah baris kode yang dihapus'; COMMENT ON COLUMN public.gh_realtime_data.changed_files IS 'Jumlah file yang diubah dalam pull request'; COMMENT ON COLUMN public.gh_realtime_data.push_size IS 'Jumlah push'; COMMENT ON COLUMN public.gh_realtime_data.push_distinct_size IS 'Jumlah push berbeda'; COMMENT ON COLUMN public.gh_realtime_data.hr IS 'Jam terjadinya event. Misalnya, jika waktu adalah 00:23, hr=00.'; COMMENT ON COLUMN public.gh_realtime_data.month IS 'Bulan terjadinya event. Misalnya, jika tanggal adalah Oktober 2015, month=2015-10.'; COMMENT ON COLUMN public.gh_realtime_data.year IS 'Tahun terjadinya event. Misalnya, jika tahun adalah 2015, year=2015.'; COMMENT ON COLUMN public.gh_realtime_data.ds IS 'Hari terjadinya event. ds=yyyy-mm-dd.'; COMMIT; -
Tuliskan data secara real-time menggunakan Flink.
Gunakan Flink untuk mengurai data SLS dan menuliskannya ke Hologres secara real-time. Pernyataan Flink berikut menyaring data: data kotor di mana ID event atau waktu event (
created_at) bernilai null dibuang, dan hanya data event terbaru yang dipertahankan.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>',--Titik akhir pribadi SLS 'accessid' = '<accesskey id>',--ID AccessKey akun Anda 'accesskey' = '<accesskey secret>',--Rahasia AccessKey akun Anda 'project' = '<project name>',--Nama proyek SLS 'logstore' = '<logstore name>'--Nama Logstore SLS 'starttime' = '2023-04-06 00:00:00',--Waktu mulai pengumpulan data 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>', --Nama database Hologres 'tablename' = '<hologres tablename>', --Nama tabel Hologres yang menerima data 'username' = '<accesskey id>', --ID AccessKey Akun Alibaba Cloud saat ini 'password' = '<accesskey secret>', --Rahasia AccessKey Akun Alibaba Cloud saat ini 'endpoint' = '<endpoint>', --Titik akhir VPC instans Hologres saat ini 'jdbcretrycount' = '1', --Jumlah percobaan ulang saat koneksi gagal 'partitionrouter' = 'true', --Menentukan apakah data ditulis ke tabel partisi 'createparttable' = 'true', --Menentukan apakah partisi dibuat otomatis 'mutatetype' = 'insertorignore' --Mode penulisan data ); 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);Untuk informasi selengkapnya tentang parameter, lihat Simple Log Service (SLS) dan Hologres.
CatatanData event mentah dari GitHub menggunakan zona waktu UTC dan tidak memiliki atribut zona waktu. Zona waktu default Hologres adalah UTC+8. Oleh karena itu, sesuaikan zona waktu saat menuliskan data dari Flink ke Hologres secara real-time. Tetapkan atribut zona waktu UTC ke data tabel sumber dalam Flink SQL. Langkah-langkahnya sebagai berikut:
Langkah 1: Buka halaman pengeditan pekerjaan
-
Masuk ke Konsol Realtime Compute for Apache Flink
-
Buka ruang kerja tujuan.
-
Temukan pekerjaan Flink SQL atau JAR Anda dan klik Edit
Langkah 2: Buka tab Detail Penyebaran
⚠️ Perubahan penting: "Konfigurasi Flink" tidak lagi menjadi bagian terpisah. Telah digabung ke dalam "Detail Penyebaran".
-
Di bagian atas halaman pengeditan pekerjaan, beralih ke tab Detail Penyebaran.
-
Gulir ke bawah ke bagian Konfigurasi Parameter.
Langkah 3: Tambahkan konfigurasi kustom
-
Klik tombol Edit di sebelah kanan Konfigurasi Parameter.
-
Pada kotak dialog yang muncul, temukan kotak teks Konfigurasi Lainnya.
-
Di kotak teks tersebut, tambahkan parameter Flink
table.local-time-zone:Asia/Shanghaisebagai pasangan kunci-nilai untuk mengatur zona waktu sistem Flink keAsia/Shanghai.
-
-
Query data.
Query data SLS yang dituliskan ke Hologres melalui Flink, lalu lakukan pengembangan data sesuai kebutuhan.
SELECT * FROM public.gh_realtime_data limit 10;Hasil query mengembalikan bidang-bidang berikut:
-
id -
actor_id -
actor_login -
repo_id -
repo_name -
org_id -
org_login -
type(jenis event, seperti PullRequestReviewEvent, CreateEvent, PushEvent, PullRequestEvent) -
created_at -
action -
iss_or_pr_id
-
Koreksi data real-time menggunakan data offline
Dalam skenario ini, data real-time mungkin hilang. Anda dapat menggunakan data offline untuk mengoreksi data real-time. Langkah-langkah berikut menunjukkan cara mengoreksi data real-time hari sebelumnya. Sesuaikan periode koreksi sesuai kebutuhan.
-
Buat tabel eksternal di Hologres untuk mendapatkan data offline 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');Untuk informasi selengkapnya tentang parameter, lihat IMPORT FOREIGN SCHEMA.
-
Buat tabel sementara untuk mengoreksi data real-time hari sebelumnya dengan data offline.
CatatanHologres V2.1.17 dan versi lebih baru mendukung Komputasi tanpa server. Untuk skenario seperti impor data offline skala besar, pekerjaan ETL besar, dan query volume tinggi pada tabel eksternal, Anda dapat menggunakan Komputasi tanpa server untuk menjalankan tugas-tugas ini. Fitur ini menggunakan sumber daya tanpa server tambahan alih-alih sumber daya instans Anda, yang meningkatkan stabilitas instans dan mengurangi kemungkinan error kehabisan memori (OOM). Anda tidak perlu menyediakan sumber daya komputasi tambahan untuk instans Anda, dan Anda hanya dikenai biaya untuk tugas yang dijalankan. Untuk informasi selengkapnya tentang Komputasi tanpa server, lihat Komputasi tanpa server. Untuk petunjuk penggunaan Komputasi tanpa server, lihat Gunakan Komputasi tanpa server.
-- Bersihkan tabel sementara yang mungkin ada DROP TABLE IF EXISTS gh_realtime_data_tmp; -- Buat tabel sementara SET hg_experimental_enable_create_table_like_properties = ON; CALL HG_CREATE_TABLE_LIKE ('gh_realtime_data_tmp', 'select * from gh_realtime_data'); -- (Opsional) Gunakan Komputasi tanpa server untuk melakukan impor data offline skala besar dan pekerjaan ETL. SET hg_computing_resource = 'serverless'; -- Masukkan data ke tabel sementara dan perbarui statistik 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; -- Atur ulang konfigurasi untuk memastikan pernyataan SQL non-esensial tidak menggunakan sumber daya tanpa server. RESET hg_computing_resource; -- Ganti tabel atomik dengan tabel anak sementara yang ada 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;
Analisis data
Anda dapat melakukan berbagai analisis pada data yang dikumpulkan. Berdasarkan rentang waktu yang dibutuhkan bisnis Anda, rancang gudang data Anda secara berlapis untuk mendukung analisis real-time, analisis offline, dan analisis real-time serta offline terpadu.
Contoh berikut menganalisis data real-time. Anda juga dapat menganalisis data untuk repositori kode atau developer tertentu.
-
Query jumlah total event publik hari ini.
SELECT count(*) FROM gh_realtime_data WHERE created_at >= date_trunc('day', now());Berikut adalah hasil sampel:
count ------ 1006 -
Query proyek paling aktif (dengan event terbanyak) dalam sehari terakhir.
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;Berikut adalah hasil sampel:
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 -
Query developer paling aktif (dengan event terbanyak) dalam sehari terakhir.
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;Berikut adalah hasil sampel:
actor_login events ------------------+------ direwolf-github 13 arm-on 10 sergii-nosachenko 9 Christoffel-T 9 yangwang201911 7 -
Query peringkat bahasa pemrograman paling populer dalam satu jam terakhir.
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;Berikut adalah hasil sampel:
language total -----------+---- JavaScript 25 C++ 15 Python 14 TypeScript 13 Java 8 PHP 8 -
Query peringkat proyek berdasarkan jumlah bintang yang diterima dalam sehari terakhir.
CatatanContoh ini tidak memperhitungkan kasus di mana pengguna membatalkan bintang pada proyek.
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;Berikut adalah hasil sampel:
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 -
Query pengguna aktif harian dan proyek hari ini.
SELECT uniq (actor_id) actor_num, uniq (repo_id) repo_num FROM gh_realtime_data WHERE created_at > date_trunc('day', now());Berikut adalah hasil sampel:
actor_num repo_num ---------+-------- 743 816