Visão geral
Empresas de jogos precisam de análises detalhadas do comportamento dos jogadores para aumentar a retenção, a monetização e o engajamento. Os data warehouses tradicionais têm custos elevados, são difíceis de dimensionar e apresentam fragmentação entre diferentes engines.
Esta solução utiliza o framework de lakehouse de código aberto Alibaba Cloud OpenLake com os seguintes componentes principais:
EMR Serverless Spark: engine Spark serverless para ETL eficiente.
Paimon: formato de armazenamento de lake para streaming e batch, com transações ACID, evolução de schema e time travel.
DLF (Data Lake Formation): gerenciamento unificado de metadados que conecta Spark, StarRocks, Flink e outras engines.
StarRocks: engine analítica de alto desempenho para consultas pontuais de alta concorrência e OLAP complexo.
Principais recursos:
Organize dados em camadas, desde logs brutos no OSS até as camadas DWD e ADS.
Execute ETL em batch com Spark.
Grave resultados em tabelas de lake do Paimon para acesso por múltiplas engines.
Consulte tabelas do lake ou internas com StarRocks para relatórios de BI e análises ad hoc.
Diagrama de arquitetura
Pré-requisitos
Antes de começar, prepare os seguintes serviços:
|
Produto |
Ação |
|
EMR Serverless |
Crie um cluster de computação Spark e anexe permissões do DLF. |
|
DLF |
Ative o serviço. Crie um catálogo Paimon na região de destino e obtenha o |
|
StarRocks |
Implante uma instância EMR Serverless StarRocks. Opcional caso utilize apenas Spark + Paimon. |
|
DataWorks |
Execute scripts de ETL Spark. |
Procedimento
Etapa 1: Configure variáveis de ambiente (execute no Notebook)
%emr_serverless_spark
DLF_CATALOG_ID = "clg-paimon-e62c8d1e8fa04ee097be4870af155296" # ← Replace with your DLF Catalog ID
REGION = "cn-hangzhou" # ← Replace with your region
print(f"DLF Catalog ID: {DLF_CATALOG_ID}")
print(f"Region: {REGION}")
Etapa 2: Inicializar sessão do Spark e verifique fontes de dados
from pyspark.sql import SparkSession
OSS_PUBLIC_BUCKET = f"emr-starrocks-benchmark-resource-{REGION}"
PROFILE_SRC_GLOB = f"oss://{OSS_PUBLIC_BUCKET}/sr_game_demo_v2/user_profile/*.parquet"
EVENT_SRC_GLOB = f"oss://{OSS_PUBLIC_BUCKET}/sr_game_demo_v2/user_event/*.parquet"
spark = (
SparkSession.builder
.appName("DLF-Paimon-Ingest-sr_game_demo_v2")
.config("spark.dlf.catalog.id", DLF_CATALOG_ID)
.config("spark.dlf.region", REGION)
.config("spark.hadoop.fs.oss.endpoint", f"oss-{REGION}-internal.aliyuncs.com")
.enableHiveSupport()
.getOrCreate()
)
# Verify that the files exist
def glob_count(path_glob: str) -> int:
from py4j.java_gateway import java_import
jvm = spark._jvm
hconf = spark.sparkContext._jsc.hadoopConfiguration()
Path = jvm.org.apache.hadoop.fs.Path
p = Path(path_glob)
fs = p.getFileSystem(hconf)
stats = fs.globStatus(p)
return 0 if stats is None else len(stats)
print("profile matched:", glob_count(PROFILE_SRC_GLOB))
print("event matched: ", glob_count(EVENT_SRC_GLOB))
Etapa 3: Ler dados brutos e gravar na camada ODS do Paimon
O Paimon gerencia automaticamente a mesclagem de arquivos, a indexação e a evolução de schema, eliminando a necessidade de particionamento manual de Parquet.
df_profile = spark.read.parquet(PROFILE_SRC_GLOB)
df_event = spark.read.parquet(EVENT_SRC_GLOB)
spark.sql("CREATE DATABASE IF NOT EXISTS game_db")
# Write the data to Paimon tables in the ODS layer
(df_profile.write
.format("paimon")
.mode("overwrite")
.saveAsTable("game_db.ods_user_profile"))
(df_event.write
.format("paimon")
.mode("overwrite")
.saveAsTable("game_db.ods_user_event"))
Etapa 4: Construir a camada de dados de detalhe DWD (ETL com Spark SQL)
4,1 Tabela de detalhes do usuário dwd_user_details
DROP TABLE IF EXISTS game_db.dwd_user_details;
CREATE TABLE game_db.dwd_user_details
USING paimon
AS
WITH p AS (
SELECT
user_id,
decode(gender, 'UTF-8') AS gender,
decode(os_version, 'UTF-8') AS os_version,
CAST(decode(current_level, 'UTF-8') AS INT) AS current_level,
decode(device_type, 'UTF-8') AS device_type,
TO_DATE(decode(last_login_date, 'UTF-8')) AS last_login_date,
decode(favorite_game_mode, 'UTF-8') AS favorite_game_mode,
decode(language_preference, 'UTF-8') AS language_preference,
decode(active_time, 'UTF-8') AS active_time,
TO_DATE(decode(registration_date, 'UTF-8')) AS registration_date,
CAST(decode(total_deaths, 'UTF-8') AS INT) AS total_deaths,
CAST(decode(game_hours, 'UTF-8') AS DOUBLE) AS game_hours,
decode(location, 'UTF-8') AS location,
decode(play_frequency, 'UTF-8') AS play_frequency,
ROW_NUMBER() OVER (
PARTITION BY user_id
ORDER BY TO_DATE(decode(last_login_date, 'UTF-8')) DESC NULLS LAST
) AS rn
FROM game_db.ods_user_profile
)
SELECT * EXCEPT (rn)
FROM p
WHERE rn = 1; -- Deduplicate: Keep only the record with the latest last_login_date for each user_id
4,2 Tabela de detalhes de eventos do usuário dwd_user_event
DROP TABLE IF EXISTS game_db.dwd_user_event;
CREATE TABLE game_db.dwd_user_event
USING paimon
AS
WITH e AS (
SELECT
user_id,
LOWER(decode(event_type, 'UTF-8')) AS event_type,
TRIM(decode(timestamp, 'UTF-8')) AS ts_str,
decode(event_details, 'UTF-8') AS event_details,
decode(location, 'UTF-8') AS event_location
FROM game_db.ods_user_event
),
e_ts AS (
SELECT *,
CASE
WHEN ts_str RLIKE '^[0-9]{13}$' THEN TO_TIMESTAMP(FROM_UNIXTIME(CAST(ts_str AS BIGINT)/1000))
WHEN ts_str RLIKE '^[0-9]{10}$' THEN TO_TIMESTAMP(FROM_UNIXTIME(CAST(ts_str AS BIGINT)))
ELSE TO_TIMESTAMP(ts_str)
END AS event_ts
FROM e
)
SELECT
e.user_id,
e.event_ts,
TO_DATE(e.event_ts) AS event_date,
e.event_type,
COALESCE(NULLIF(e.event_location,''), d.location) AS location,
-- Extract amount from JSON
COALESCE(CAST(get_json_object(event_details, '$.amount') AS DOUBLE), 0.0) AS amount,
d.gender, d.device_type
FROM e_ts e
LEFT JOIN game_db.dwd_user_details d ON e.user_id = d.user_id
WHERE e.event_ts IS NOT NULL;
Etapa 5: Construir a camada de dados de aplicação ADS (agregação de métricas)
5,1 Tabela de taxa de retenção diária ads_retention_daily
DROP TABLE IF EXISTS game_db.ads_retention_daily;
CREATE TABLE game_db.ads_retention_daily
USING paimon
AS
WITH dau AS (
SELECT event_date AS dt, user_id
FROM game_db.dwd_user_event
GROUP BY event_date, user_id
)
SELECT
base.dt,
base.dau,
COALESCE(d1.d1_retained, 0) AS d1_retained,
ROUND(COALESCE(d1.d1_retained,0) * 1.0 / base.dau, 4) AS d1_retention_rate,
COALESCE(d7.d7_retained, 0) AS d7_retained,
ROUND(COALESCE(d7.d7_retained,0) * 1.0 / base.dau, 4) AS d7_retention_rate
FROM (
SELECT dt, COUNT(DISTINCT user_id) AS dau FROM dau GROUP BY dt
) base
LEFT JOIN (
SELECT a.dt, COUNT(DISTINCT a.user_id) AS d1_retained
FROM dau a JOIN dau b ON a.user_id = b.user_id AND b.dt = DATE_ADD(a.dt, 1)
GROUP BY a.dt
) d1 ON base.dt = d1.dt
LEFT JOIN (
SELECT a.dt, COUNT(DISTINCT a.user_id) AS d7_retained
FROM dau a JOIN dau b ON a.user_id = b.user_id AND b.dt = DATE_ADD(a.dt, 7)
GROUP BY a.dt
) d7 ON base.dt = d7.dt;
5,2 Outras tabelas ADS (consulte o código original)
ads_purchase_trends_daily: GMV diário e usuários pagantes.ads_device_preference_daily: distribuição e participação de dispositivos.ads_region_distribution_daily: distribuição de DAU por província e município.
Etapa 6: Conectar ao StarRocks
Crie um nó StarRocks para consultar tabelas do lake diretamente.
-- Query lake tables directly
SELECT * FROM paimon_catalog.game_db.ads_retention_daily LIMIT 10;
A consulta retorna as colunas dt, dau, d1_retained, d1_retention_rate, d7_retained e d7_retention_rate, exibindo dados de retenção por data em ordem decrescente.
Etapa 7: Visualizar dados no Quick BI
No console do Quick BI, crie uma fonte de dados StarRocks.
-
Crie conjuntos de dados com o seguinte SQL:
SQL do conjunto de dados 1:
SELECT * FROM game_db.ADS_MV_USER_RETENTION;SQL do conjunto de dados 2:
SELECT * FROM game_db.ADS_MV_USER_GEOGRAPHIC_DISTRIBUTION;SQL do conjunto de dados 3:
SELECT * FROM game_db.ADS_MV_USER_DEVICE_PREFERENCE;SQL do conjunto de dados 4:
SELECT * FROM game_db.ADS_MV_USER_PURCHASE_TRENDS;
Crie gráficos de linhas, mapas e painéis para monitorar as principais métricas.
Benefícios
|
Dimensão |
Abordagem tradicional |
Esta solução |
|
Custo de armazenamento |
Alto custo com HDFS |
OSS de baixo custo + mesclagem automática de arquivos pequenos no Paimon |
|
Elasticidade de computação |
Exige clusters sempre ativos |
EMR Serverless com pagamento conforme o uso |
|
Consistência de dados |
Gerenciamento manual de partições e versões |
ACID do Paimon + time travel |
|
Colaboração multi-engine |
Silos de dados |
Metadados unificados no DLF, compartilhados por Spark, Flink e StarRocks |
|
Eficiência de desenvolvimento |
Dependências complexas de agendamento |
ETL ponta a ponta e modelagem SQL no Notebook |