Topik ini menjelaskan cara melakukan deduplikasi unique visitor (UV) waktu nyata yang akurat menggunakan Hologres dan Flink.
Prasyarat
-
Sebuah instans Hologres aktif telah terhubung ke alat pengembangan. Topik ini menggunakan HoloWeb sebagai contoh. Untuk informasi lebih lanjut, lihat Hubungkan ke dan kueri Hologres menggunakan HoloWeb.
-
Kluster Flink telah disiapkan. Anda dapat menggunakan Realtime Compute for Apache Flink atau Apache Flink.
Informasi latar belakang
Hologres sangat terintegrasi dengan Flink, mendukung penulisan data waktu nyata ber-throughput tinggi dengan visibilitas langsung. Hologres juga mendukung Flink SQL untuk join tabel dimensi dan pengembangan berbasis event menggunakan sumber change data capture (CDC). Kombinasi kuat ini menjadikannya ideal untuk deduplikasi UV waktu nyata. Diagram berikut menunjukkan arsitektur solusi.
-
Flink berlangganan data waktu nyata dari berbagai sumber data, seperti log dari Kafka.
-
Flink memproses data dengan mengonversi aliran data menjadi tabel, melakukan join dengan tabel dimensi Hologres, dan menulis hasilnya ke Hologres secara waktu nyata.
-
Hologres memproses data yang ditulis dari Flink secara waktu nyata.
-
Aplikasi hulu, seperti DataService Studio atau Quick BI, mengonsumsi hasil kueri akhir.
Alur kerja penghitungan UV waktu nyata
Integrasi kuat antara Flink dan Hologres, dikombinasikan dengan dukungan native untuk tipe data roaring bitmap di Hologres, memungkinkan penghitungan UV waktu nyata dan deduplikasi tag pengguna. Diagram berikut menunjukkan alur kerja detail.
-
Flink berlangganan data pengguna secara waktu nyata dari sumber data seperti Kafka atau Redis dan mengonversi aliran data menjadi tabel sumber.
-
Di Hologres, buat tabel pemetaan pengguna untuk menyimpan ID pengguna historis (UID) dan UID auto-increment 32-bit yang sesuai.
CatatanID pengguna dari sistem bisnis atau titik pelacakan umumnya berupa string atau integer panjang. Namun, tipe data roaring bitmap mengharuskan UID berupa integer 32-bit. Untuk performa terbaik, integer tersebut sebaiknya sedens mungkin (berurutan). Tabel pemetaan menggunakan tipe
SERIALHologres, yaitu integer 32-bit auto-increment, untuk secara otomatis mempertahankan pemetaan stabil dari UID asli ke UID integer 32-bit. -
Di Flink, gunakan tabel pemetaan pengguna Hologres sebagai tabel dimensi Flink. Gunakan fitur
insertIfNotExistsdari tabel dimensi bersama dengan auto-increment field untuk memetakan UID secara efisien. Lakukan join tabel sumber dengan tabel dimensi dan konversi hasilnya menjadi DataStream. -
Buat tabel hasil agregasi di Hologres. Flink memproses data yang telah di-join dalam jendela waktu dan menerapkan fungsi roaring bitmap berdasarkan dimensi kueri yang diinginkan.
-
Untuk mengkueri data, pilih dari tabel hasil agregasi berdasarkan kondisi kueri Anda. Lakukan operasi
ORpada bidang roaring bitmap yang relevan dan hitung kardinalitas untuk mendapatkan jumlah pengguna akhir.
Pendekatan ini menyediakan data UV dan page view (PV) waktu nyata dengan granularitas detail halus. Pendekatan ini memungkinkan Anda menyesuaikan jendela statistik minimum, seperti UV dalam 5 menit terakhir, untuk mengaktifkan Pemantauan waktu nyata pada tampilan BI seperti layar besar. Dibandingkan dengan deduplikasi per hari, minggu, atau bulan, metode ini lebih cocok untuk analisis detail halus selama event tertentu. Anda juga dapat memperoleh hasil untuk rentang waktu yang lebih panjang melalui agregasi sederhana. Namun, jika Anda mengagregasi data dengan granularitas detail tetapi mengkueri tanpa filter atau dimensi agregasi yang sesuai, Anda mungkin memicu operasi agregasi tambahan saat kueri, yang dapat menurunkan performa.
Solusi ini memiliki pipa data yang sederhana, memungkinkan komputasi fleksibel di semua dimensi, dan menggunakan satu bitmap untuk penyimpanan, sehingga menghindari masalah ledakan penyimpanan. Solusi ini juga menjamin pembaruan waktu nyata untuk menciptakan gudang data analitik multidimensi yang lebih responsif, fleksibel, dan andal.
Prosedur
-
Buat tabel dasar di Hologres
-
Buat tabel pemetaan pengguna
Di Hologres, jalankan pernyataan berikut untuk membuat tabel pemetaan pengguna bernama
uid_mapping. Tabel ini memetakan UID ke integer 32-bit. Jika UID asli Anda sudah berupa integer 32-bit, Anda dapat melewati langkah ini.-
ID pengguna dari sistem bisnis atau titik pelacakan umumnya berupa string atau integer panjang. Oleh karena itu, Anda perlu membuat tabel
uid_mapping. Tipe data roaring bitmap mengharuskan ID pengguna berupa integer 32-bit yang sedens mungkin (sebaiknya berurutan). Tabel pemetaan menggunakan tipeSERIALHologres, yaitu integer 32-bit auto-increment, untuk secara otomatis mengelola dan mempertahankan pemetaan stabil. -
Untuk meningkatkan permintaan per detik (QPS) join tabel dimensi Flink, atur tabel ini di Hologres sebagai Tabel berorientasi baris.
-
Anda harus mengaktifkan parameter GUC yang sesuai untuk menggunakan Mesin eksekusi yang dioptimalkan saat menulis data ke tabel yang berisi
kolom auto-increment. Untuk informasi lebih lanjut, lihat Percepat eksekusi SQL menggunakan Fixed Plan.
-- Aktifkan parameter GUC untuk mendukung penulisan Fixed Plan pada tabel yang berisi kolom bertipe SERIAL. alter database <dbname> set hg_experimental_enable_fixed_dispatcher_autofill_series=on; alter database <dbname> set hg_experimental_enable_fixed_dispatcher_for_multi_values=on; BEGIN; CREATE TABLE public.uid_mapping ( uid text NOT NULL, uid_int32 serial, PRIMARY KEY (uid) ); -- Atur uid sebagai clustering_key dan distribution_key untuk menemukan nilai int32-nya dengan cepat. CALL set_table_property('public.uid_mapping', 'clustering_key', 'uid'); CALL set_table_property('public.uid_mapping', 'distribution_key', 'uid'); CALL set_table_property('public.uid_mapping', 'orientation', 'row'); COMMIT; -
-
Buat tabel hasil agregasi
Buat tabel bernama
dws_appsebagai tabel hasil agregasi untuk menyimpan hasil yang diagregasi berdasarkan dimensi dasar.Sebelum menggunakan fungsi roaring bitmap, Anda harus membuat ekstensi roaringbitmap. Instans Hologres Anda harus versi V0.10 atau lebih baru.
CREATE EXTENSION IF NOT EXISTS roaringbitmap;Dibandingkan dengan tabel hasil offline, tabel ini mencakup bidang timestamp untuk memungkinkan statistik berdasarkan periode jendela Flink. Pernyataan DDL berikut mendefinisikan tabel hasil.
BEGIN; CREATE TABLE dws_app( country text, prov text, city text, ymd text NOT NULL, -- Bidang tanggal timetz TIMESTAMPTZ, -- Timestamp statistik, memungkinkan perhitungan statistik berdasarkan periode jendela Flink. uid32_bitmap roaringbitmap, -- Gunakan roaring bitmap untuk mencatat UV. PRIMARY KEY (country, prov, city, ymd, timetz)-- Gunakan dimensi kueri dan waktu sebagai kunci primer untuk mencegah penyisipan data duplikat. ); CALL set_table_property('public.dws_app', 'orientation', 'column'); -- Atur bidang tanggal sebagai clustering key dan event_time_column untuk filtering efisien. CALL set_table_property('public.dws_app', 'clustering_key', 'ymd'); CALL set_table_property('public.dws_app', 'event_time_column', 'ymd'); -- Atur bidang GROUP BY sebagai distribution key. CALL set_table_property('public.dws_app', 'distribution_key', 'country,prov,city'); COMMIT;
-
-
Gunakan Flink untuk membaca data secara waktu nyata dan memperbarui tabel hasil agregasi
Untuk kode sumber lengkap contoh Flink, lihat contoh alibabacloud-hologres-connectors. Langkah-langkah berikut menjelaskan operasi di Flink.
-
Baca data dari sumber data sebagai DataStream dan konversi DataStream menjadi Table
Di Flink, baca data dari sumber data, seperti file CSV, Kafka, atau Redis.
// File CSV digunakan sebagai sumber data dalam contoh ini. Anda juga dapat menggunakan sumber data lain seperti Kafka atau Redis. DataStreamSource odsStream = env.createInput(csvInput, typeInfo); // Bidang proctime harus ditambahkan untuk join dengan tabel dimensi. Table odsTable = tableEnv.fromDataStream( odsStream, $("uid"), $("country"), $("prov"), $("city"), $("ymd"), $("proctime").proctime()); // Daftarkan tabel di lingkungan katalog. tableEnv.createTemporaryView("odsTable", odsTable); -
Lakukan join tabel sumber dengan tabel dimensi Hologres (uid_mapping)
Saat membuat tabel dimensi Hologres di Flink, gunakan parameter
insertIfNotExistsuntuk secara otomatis menyisipkan data yang tidak ada. Bidanguid_int32dihasilkan secara otomatis menggunakan tipeSERIALHologres. Lakukan join tabel sumber Flink dengan tabel dimensi Hologres. Kode berikut memberikan contohnya.-- Buat tabel dimensi Hologres. 'insertIfNotExists' menunjukkan bahwa jika data tidak ditemukan, data tersebut akan disisipkan secara otomatis. String createUidMappingTable = String.format( "create table uid_mapping_dim(" + " uid string," + " uid_int32 INT" + ") with (" + " 'connector'='hologres'," + " 'dbname' = '%s'," // Nama database Hologres. + " 'tablename' = '%s',"// Nama tabel Hologres. + " 'username' = '%s'," // ID AccessKey akun Anda. + " 'password' = '%s'," // Rahasia AccessKey akun Anda. + " 'endpoint' = '%s'," // Titik akhir instans Hologres. + " 'insertifnotexists'='true'" + ")", database, dimTableName, username, password, endpoint); tableEnv.executeSql(createUidMappingTable); -- Lakukan join tabel sumber dengan tabel dimensi. String odsJoinDim = "SELECT ods.country, ods.prov, ods.city, ods.ymd, dim.uid_int32" + " FROM odsTable AS ods JOIN uid_mapping_dim FOR SYSTEM_TIME AS OF ods.proctime AS dim" + " ON ods.uid = dim.uid"; Table joinRes = tableEnv.sqlQuery(odsJoinDim); -
Konversi hasil join menjadi DataStream
Proses data menggunakan jendela waktu Flink dan gunakan roaring bitmap untuk mendeduplikasi metrik. Kode berikut memberikan contohnya.
DataStream<Tuple6<String, String, String, String, Timestamp, byte[]>> processedSource = source -- Filter dimensi yang memerlukan statistik. Dalam contoh ini, dimensinya adalah country, prov, city, dan ymd. .keyBy(0, 1, 2, 3) -- Jendela waktu tumbling. Karena file CSV digunakan untuk mensimulasikan aliran input, ProcessingTime digunakan. Di lingkungan produksi, Anda dapat menggunakan EventTime. .window(TumblingProcessingTimeWindows.of(Time.minutes(5))) -- Pemicu yang memungkinkan Anda memperoleh hasil agregat sebelum jendela ditutup. .trigger(ContinuousProcessingTimeTrigger.of(Time.minutes(1))) .aggregate( -- Fungsi agregat mengagregasi data berdasarkan dimensi yang dipilih di keyBy. new AggregateFunction< Tuple5<String, String, String, String, Integer>, RoaringBitmap, RoaringBitmap>() { @Override public RoaringBitmap createAccumulator() { return new RoaringBitmap(); } @Override public RoaringBitmap add( Tuple5<String, String, String, String, Integer> in, RoaringBitmap acc) { -- Tambahkan UID 32-bit ke roaring bitmap untuk deduplikasi. acc.add(in.f4); return acc; } @Override public RoaringBitmap getResult(RoaringBitmap acc) { return acc; } @Override public RoaringBitmap merge( RoaringBitmap acc1, RoaringBitmap acc2) { return RoaringBitmap.or(acc1, acc2); } }, -- Fungsi jendela mengeluarkan hasil agregat. new WindowFunction< RoaringBitmap, Tuple6<String, String, String, String, Timestamp, byte[]>, Tuple, TimeWindow>() { @Override public void apply( Tuple keys, TimeWindow timeWindow, Iterable<RoaringBitmap> iterable, Collector< Tuple6<String, String, String, String, Timestamp, byte[]>> out) throws Exception { RoaringBitmap result = iterable.iterator().next(); // Optimalkan roaring bitmap. result.runOptimize(); // Konversi roaring bitmap menjadi array byte untuk menyimpannya di Hologres. byte[] byteArray = new byte[result.serializedSizeInBytes()]; result.serialize(ByteBuffer.wrap(byteArray)); // Timestamp (Tuple6.f4) menandai akhir jendela waktu, yang menentukan granularitas statistik. out.collect( new Tuple6<>( keys.getField(0), keys.getField(1), keys.getField(2), keys.getField(3), new Timestamp( timeWindow.getEnd() / 1000 * 1000), byteArray)); } }); -
Tulis data ke tabel hasil agregasi Hologres
Tulis data yang telah dideduplikasi oleh Flink ke tabel hasil Hologres dws_app. Perhatikan bahwa tipe
roaringbitmapdi Hologres berkorespondensi dengan tipebyte arraydi Flink. Kode berikut memberikan contohnya di Flink.-- Konversi hasil perhitungan menjadi tabel. Table resTable = tableEnv.fromDataStream( processedSource, $("country"), $("prov"), $("city"), $("ymd"), $("timest"), $("uid32_bitmap")); -- Buat tabel sink Hologres. Tipe roaringbitmap di Hologres berkorespondensi dengan tipe BYTES dalam definisi tabel Flink. String createHologresTable = String.format( "create table sink(" + " country string," + " prov string," + " city string," + " ymd string," + " timetz timestamp," + " uid32_bitmap BYTES" + ") with (" + " 'connector'='hologres'," + " 'dbname' = '%s'," + " 'tablename' = '%s'," + " 'username' = '%s'," + " 'password' = '%s'," + " 'endpoint' = '%s'," + " 'connectionSize' = '%s'," + " 'mutatetype' = 'insertOrReplace'" + ")", database, dwsTableName, username, password, endpoint, connectionSize); tableEnv.executeSql(createHologresTable); -- Tulis hasil perhitungan ke tabel dws_app. tableEnv.executeSql("insert into sink select * from " + resTable);
-
-
Kueri data
Di Hologres, hitung UV dari tabel hasil agregasi (
dws_app). Agregasi data berdasarkan dimensi kueri Anda dan hitung kardinalitas bitmap untuk menghitung pengguna yang sesuai dengan kondisiGROUP BY.-
Contoh 1: Kueri jumlah UV untuk setiap kota pada hari tertentu.
-- Sebelum menjalankan operasi RB_AGG, Anda dapat menonaktifkan switch agregasi tiga tahap untuk performa lebih baik. Langkah ini opsional karena switch tersebut dinonaktifkan secara default. set hg_experimental_enable_force_three_stage_agg=off; SELECT country ,prov ,city ,RB_CARDINALITY(RB_OR_AGG(uid32_bitmap)) AS uv FROM dws_app WHERE ymd = '20210329' GROUP BY country ,prov ,city ; -
Contoh 2: Kueri jumlah UV dan PV untuk setiap provinsi dalam periode waktu tertentu.
-- Sebelum menjalankan operasi RB_AGG, Anda dapat menonaktifkan switch agregasi tiga tahap untuk performa lebih baik. Langkah ini opsional karena switch tersebut dinonaktifkan secara default. set hg_experimental_enable_force_three_stage_agg=off; SELECT country ,prov ,RB_CARDINALITY(RB_OR_AGG(uid32_bitmap)) AS uv ,SUM(pv) AS pv FROM dws_app WHERE timetz > '2021-04-19 18:00:00+08' and timetz < '2021-04-19 19:00:00+08' GROUP BY country ,prov ;
-
-
Visualisasikan hasil
Setelah menghitung UV dan PV, Anda biasanya menggunakan alat BI untuk visualisasi. Karena kueri memerlukan fungsi agregat RB_CARDINALITY dan RB_OR_AGG, Anda memerlukan alat BI yang mendukung fungsi agregat kustom. Alat BI umum yang memiliki kemampuan ini termasuk Apache Superset dan Tableau.
-
Apache Superset
-
Hubungkan Apache Superset ke Hologres. Untuk informasi lebih lanjut, lihat Hubungkan Apache Superset ke Hologres.
-
Atur tabel dws_app sebagai set data. Di Apache Superset, klik Add Dataset. Di kotak dialog yang muncul, atur DATASOURCE ke
postgresql holo_rb_demo, SCHEMA kepublic, dan TABLE kedws_app. -
Di set data, buat metrik bernama UV menggunakan ekspresi berikut. Di halaman pengeditan set data, pilih tab METRICS dan klik + ADD ITEM untuk menambahkan metrik. Atur Metric ke
countdan SQL Expression keCOUNT(*). Atur Metric keuvdan SQL Expression keRB_CARDINALITY(RB_OR_AGG(uid32_bitmap)). Klik SAVE.RB_CARDINALITY(RB_OR_AGG(uid32_bitmap))Anda sekarang dapat menjelajahi data.
-
(Opsional) Buat Dasbor.
Untuk informasi lebih lanjut tentang cara membuat Dasbor, lihat Create a Dashboard.
-
-
Tableau
-
Hubungkan Tableau ke Hologres. Untuk informasi lebih lanjut, lihat Hubungkan Tableau ke Hologres.
Anda dapat menggunakan fungsi transmisi langsung Tableau untuk menjalankan fungsi kustom. Untuk informasi lebih lanjut, lihat Pass-Through Functions (RAWSQL).
-
Di Tableau, buat bidang terhitung bernama UV dan masukkan rumus berikut.
RAWSQLAGG_INT("RB_CARDINALITY(RB_OR_AGG(%1))", [Uid32 Bitmap])Anda sekarang dapat menjelajahi data.
-
(Opsional) Buat Dasbor.
Untuk informasi lebih lanjut tentang cara membuat Dasbor, lihat Create a Dashboard.
-
-