Todos os produtos
Search
Central de documentação

OpenLake:Open source lakehouse

Última atualização: Jun 29, 2026

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 Catalog ID.

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

  1. No console do Quick BI, crie uma fonte de dados StarRocks.

  2. Crie conjuntos de dados com o seguinte SQL:

    1. SQL do conjunto de dados 1: SELECT * FROM game_db.ADS_MV_USER_RETENTION;

    2. SQL do conjunto de dados 2: SELECT * FROM game_db.ADS_MV_USER_GEOGRAPHIC_DISTRIBUTION;

    3. SQL do conjunto de dados 3: SELECT * FROM game_db.ADS_MV_USER_DEVICE_PREFERENCE;

    4. SQL do conjunto de dados 4: SELECT * FROM game_db.ADS_MV_USER_PURCHASE_TRENDS;

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