wuchong commented on code in PR #2669: URL: https://github.com/apache/fluss/pull/2669#discussion_r3818158651
########## website/docs/quickstart/user_profile.md: ########## @@ -0,0 +1,290 @@ +--- +title: Real-Time User Profile +sidebar_position: 4 +--- + +# Real-Time User Profile + +This tutorial demonstrates how to build a real-time user profiling system using three core Apache Fluss features: the **Auto-Increment Column**, the **Aggregation Merge Engine**, and the built-in **RoaringBitmap SQL functions**. You will learn how to automatically map high-cardinality email identifiers to compact integer UIDs, accumulate click metrics, and count unique visitors — all directly in the storage layer, keeping the Flink job entirely stateless. + +## How the System Works + +### Core Concepts + +- **Identity Mapping**: Incoming email strings are automatically mapped to compact `INT` UIDs using Fluss's auto-increment column — no manual ID management required. +- **Storage-Level Aggregation**: Click counts are summed and unique visitor bitmaps are OR-ed directly inside the Fluss TabletServers via the Aggregation Merge Engine. +- **Built-in Bitmap Functions**: `rb_build_agg` and `rb_cardinality` are registered natively in FlussCatalog — no external JAR or `CREATE TEMPORARY FUNCTION` required. + +### Data Flow + +1. **Ingestion**: Raw click events arrive with an email address and a click count. +2. **Mapping**: A Flink lookup join against `user_dict` resolves the email to a UID. If the email is new, the `insert-if-not-exists` hint instructs Fluss to generate a new UID automatically. +3. **Aggregation**: The resolved UID is written to `user_profiles`. The Aggregation Merge Engine sums `total_clicks` and OR-s the `unique_visitors` bitmap at the storage layer — no windowing or Flink state required. + +## Prerequisites + +Before proceeding, ensure that [Docker](https://docs.docker.com/engine/install/) and the [Docker Compose plugin](https://docs.docker.com/compose/install/linux/) are installed on your machine. + +## Environment Setup + +1. Create a working directory and navigate into it. + ```shell + mkdir fluss-user-profile + cd fluss-user-profile + ``` + +2. Create a `docker-compose.yml` file with the following content: + ```yaml + services: + coordinator-server: + image: apache/fluss:$FLUSS_DOCKER_VERSION$ + command: coordinatorServer + depends_on: + - zookeeper + environment: + - | + FLUSS_PROPERTIES= + zookeeper.address: zookeeper:2181 + bind.listeners: FLUSS://coordinator-server:9123 + remote.data.dir: /tmp/fluss/remote + volumes: + - fluss-remote-data:/tmp/fluss/remote + tablet-server: + image: apache/fluss:$FLUSS_DOCKER_VERSION$ + command: tabletServer + depends_on: + - coordinator-server + environment: + - | + FLUSS_PROPERTIES= + zookeeper.address: zookeeper:2181 + bind.listeners: FLUSS://tablet-server:9123 + data.dir: /tmp/fluss/data + remote.data.dir: /tmp/fluss/remote + volumes: + - fluss-remote-data:/tmp/fluss/remote + zookeeper: + restart: always + image: zookeeper:3.9.2 + jobmanager: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + ports: + - "8083:8081" + command: jobmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.port: 8081 + volumes: + - fluss-remote-data:/tmp/fluss/remote + taskmanager: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + depends_on: + - jobmanager + command: taskmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + taskmanager.numberOfTaskSlots: 2 + volumes: + - fluss-remote-data:/tmp/fluss/remote + sql-client: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + command: ["/opt/sql-client/sql-client"] + depends_on: + - jobmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.port: 8081 + volumes: + - fluss-remote-data:/tmp/fluss/remote + + volumes: + fluss-remote-data: + ``` + +3. Start all services. + ```shell + docker compose up -d + ``` + +4. Confirm all containers are running. + ```shell + docker compose ps + ``` + You should see `coordinator-server`, `tablet-server`, `zookeeper`, `jobmanager`, `taskmanager`, and `sql-client` all in the `running` state. + +:::note +All the following commands involving `docker compose` should be executed in the working directory that contains the `docker-compose.yml` file. +::: + +## Enter the SQL Client + +Use the following command to enter the Flink SQL Client: + +```shell +docker compose run --entrypoint bash sql-client -c " +\${FLINK_HOME}/bin/sql-client.sh \ + -Drest.address=jobmanager \ + -Drest.port=8081 \ + -i /opt/sql-client/sql/sql-client.sql +" +``` + +## Step 1: Create the Fluss Catalog Review Comment: All steps should use H3 to make them under the `Enter the SQL Client` section. ########## website/docs/quickstart/user_profile.md: ########## @@ -0,0 +1,290 @@ +--- +title: Real-Time User Profile +sidebar_position: 4 +--- + +# Real-Time User Profile + +This tutorial demonstrates how to build a real-time user profiling system using three core Apache Fluss features: the **Auto-Increment Column**, the **Aggregation Merge Engine**, and the built-in **RoaringBitmap SQL functions**. You will learn how to automatically map high-cardinality email identifiers to compact integer UIDs, accumulate click metrics, and count unique visitors — all directly in the storage layer, keeping the Flink job entirely stateless. + +## How the System Works + +### Core Concepts + +- **Identity Mapping**: Incoming email strings are automatically mapped to compact `INT` UIDs using Fluss's auto-increment column — no manual ID management required. +- **Storage-Level Aggregation**: Click counts are summed and unique visitor bitmaps are OR-ed directly inside the Fluss TabletServers via the Aggregation Merge Engine. +- **Built-in Bitmap Functions**: `rb_build_agg` and `rb_cardinality` are registered natively in FlussCatalog — no external JAR or `CREATE TEMPORARY FUNCTION` required. + +### Data Flow + +1. **Ingestion**: Raw click events arrive with an email address and a click count. +2. **Mapping**: A Flink lookup join against `user_dict` resolves the email to a UID. If the email is new, the `insert-if-not-exists` hint instructs Fluss to generate a new UID automatically. +3. **Aggregation**: The resolved UID is written to `user_profiles`. The Aggregation Merge Engine sums `total_clicks` and OR-s the `unique_visitors` bitmap at the storage layer — no windowing or Flink state required. + +## Prerequisites + +Before proceeding, ensure that [Docker](https://docs.docker.com/engine/install/) and the [Docker Compose plugin](https://docs.docker.com/compose/install/linux/) are installed on your machine. + +## Environment Setup + +1. Create a working directory and navigate into it. + ```shell + mkdir fluss-user-profile + cd fluss-user-profile + ``` + +2. Create a `docker-compose.yml` file with the following content: + ```yaml + services: + coordinator-server: + image: apache/fluss:$FLUSS_DOCKER_VERSION$ + command: coordinatorServer + depends_on: + - zookeeper + environment: + - | + FLUSS_PROPERTIES= + zookeeper.address: zookeeper:2181 + bind.listeners: FLUSS://coordinator-server:9123 + remote.data.dir: /tmp/fluss/remote + volumes: + - fluss-remote-data:/tmp/fluss/remote + tablet-server: + image: apache/fluss:$FLUSS_DOCKER_VERSION$ + command: tabletServer + depends_on: + - coordinator-server + environment: + - | + FLUSS_PROPERTIES= + zookeeper.address: zookeeper:2181 + bind.listeners: FLUSS://tablet-server:9123 + data.dir: /tmp/fluss/data + remote.data.dir: /tmp/fluss/remote + volumes: + - fluss-remote-data:/tmp/fluss/remote + zookeeper: + restart: always + image: zookeeper:3.9.2 + jobmanager: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + ports: + - "8083:8081" + command: jobmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.port: 8081 + volumes: + - fluss-remote-data:/tmp/fluss/remote + taskmanager: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + depends_on: + - jobmanager + command: taskmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + taskmanager.numberOfTaskSlots: 2 + volumes: + - fluss-remote-data:/tmp/fluss/remote + sql-client: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + command: ["/opt/sql-client/sql-client"] + depends_on: + - jobmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.port: 8081 + volumes: + - fluss-remote-data:/tmp/fluss/remote + + volumes: + fluss-remote-data: + ``` + +3. Start all services. + ```shell + docker compose up -d + ``` + +4. Confirm all containers are running. + ```shell + docker compose ps + ``` + You should see `coordinator-server`, `tablet-server`, `zookeeper`, `jobmanager`, `taskmanager`, and `sql-client` all in the `running` state. + +:::note +All the following commands involving `docker compose` should be executed in the working directory that contains the `docker-compose.yml` file. +::: + +## Enter the SQL Client + +Use the following command to enter the Flink SQL Client: + +```shell +docker compose run --entrypoint bash sql-client -c " Review Comment: Can we simplify this to `docker compose run sql-client`, as in the existing Flink quickstart? The service already sets `command: ["/opt/sql-client/sql-client"]`, and the required REST settings are provided through `FLINK_PROPERTIES`, so the custom entrypoint and duplicated arguments are unnecessary. Since this tutorial creates its own source table, we may also want to point the service command at `/opt/flink/bin/sql-client.sh` to avoid preloading the unrelated demo tables from `sql-client.sql`. ########## website/docs/quickstart/user_profile.md: ########## @@ -0,0 +1,290 @@ +--- +title: Real-Time User Profile +sidebar_position: 4 +--- + +# Real-Time User Profile + +This tutorial demonstrates how to build a real-time user profiling system using three core Apache Fluss features: the **Auto-Increment Column**, the **Aggregation Merge Engine**, and the built-in **RoaringBitmap SQL functions**. You will learn how to automatically map high-cardinality email identifiers to compact integer UIDs, accumulate click metrics, and count unique visitors — all directly in the storage layer, keeping the Flink job entirely stateless. + +## How the System Works + +### Core Concepts + +- **Identity Mapping**: Incoming email strings are automatically mapped to compact `INT` UIDs using Fluss's auto-increment column — no manual ID management required. +- **Storage-Level Aggregation**: Click counts are summed and unique visitor bitmaps are OR-ed directly inside the Fluss TabletServers via the Aggregation Merge Engine. +- **Built-in Bitmap Functions**: `rb_build_agg` and `rb_cardinality` are registered natively in FlussCatalog — no external JAR or `CREATE TEMPORARY FUNCTION` required. + +### Data Flow + +1. **Ingestion**: Raw click events arrive with an email address and a click count. +2. **Mapping**: A Flink lookup join against `user_dict` resolves the email to a UID. If the email is new, the `insert-if-not-exists` hint instructs Fluss to generate a new UID automatically. +3. **Aggregation**: The resolved UID is written to `user_profiles`. The Aggregation Merge Engine sums `total_clicks` and OR-s the `unique_visitors` bitmap at the storage layer — no windowing or Flink state required. + +## Prerequisites + +Before proceeding, ensure that [Docker](https://docs.docker.com/engine/install/) and the [Docker Compose plugin](https://docs.docker.com/compose/install/linux/) are installed on your machine. + +## Environment Setup + +1. Create a working directory and navigate into it. + ```shell + mkdir fluss-user-profile + cd fluss-user-profile + ``` + +2. Create a `docker-compose.yml` file with the following content: + ```yaml + services: + coordinator-server: + image: apache/fluss:$FLUSS_DOCKER_VERSION$ + command: coordinatorServer + depends_on: + - zookeeper + environment: + - | + FLUSS_PROPERTIES= + zookeeper.address: zookeeper:2181 + bind.listeners: FLUSS://coordinator-server:9123 + remote.data.dir: /tmp/fluss/remote + volumes: + - fluss-remote-data:/tmp/fluss/remote + tablet-server: + image: apache/fluss:$FLUSS_DOCKER_VERSION$ + command: tabletServer + depends_on: + - coordinator-server + environment: + - | + FLUSS_PROPERTIES= + zookeeper.address: zookeeper:2181 + bind.listeners: FLUSS://tablet-server:9123 + data.dir: /tmp/fluss/data + remote.data.dir: /tmp/fluss/remote + volumes: + - fluss-remote-data:/tmp/fluss/remote + zookeeper: + restart: always + image: zookeeper:3.9.2 + jobmanager: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + ports: + - "8083:8081" + command: jobmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.port: 8081 + volumes: + - fluss-remote-data:/tmp/fluss/remote + taskmanager: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + depends_on: + - jobmanager + command: taskmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + taskmanager.numberOfTaskSlots: 2 + volumes: + - fluss-remote-data:/tmp/fluss/remote + sql-client: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + command: ["/opt/sql-client/sql-client"] + depends_on: + - jobmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.port: 8081 + volumes: + - fluss-remote-data:/tmp/fluss/remote + + volumes: + fluss-remote-data: + ``` + +3. Start all services. + ```shell + docker compose up -d + ``` + +4. Confirm all containers are running. + ```shell + docker compose ps + ``` + You should see `coordinator-server`, `tablet-server`, `zookeeper`, `jobmanager`, `taskmanager`, and `sql-client` all in the `running` state. + +:::note +All the following commands involving `docker compose` should be executed in the working directory that contains the `docker-compose.yml` file. +::: + +## Enter the SQL Client + +Use the following command to enter the Flink SQL Client: + +```shell +docker compose run --entrypoint bash sql-client -c " +\${FLINK_HOME}/bin/sql-client.sh \ + -Drest.address=jobmanager \ + -Drest.port=8081 \ + -i /opt/sql-client/sql/sql-client.sql +" +``` + +## Step 1: Create the Fluss Catalog + +Run these statements one by one in the SQL Client. + +:::tip +Run SQL statements one by one to avoid errors. +::: + +```sql +CREATE CATALOG fluss_catalog WITH ( + 'type' = 'fluss', + 'bootstrap.servers' = 'coordinator-server:9123' +); +``` + +```sql +USE CATALOG fluss_catalog; +``` + +:::note +Once you switch to the Fluss catalog, all RoaringBitmap SQL functions (`rb_build_agg`, `rb_cardinality`, `rb_or_agg`, and others) are available immediately — no `CREATE TEMPORARY FUNCTION` statement is needed. +::: + +## Step 2: Create the User Dictionary Table + +Create the `user_dict` table to map email addresses to integer UIDs. The `auto-increment.fields` property instructs Fluss to automatically assign a unique `INT` UID for every new email it receives. + +```sql +CREATE TABLE user_dict ( + email STRING, + uid INT, + PRIMARY KEY (email) NOT ENFORCED +) WITH ( + 'auto-increment.fields' = 'uid' +); +``` + +## Step 3: Create the Aggregated Profile Table + +Create the `user_profiles` table using the **Aggregation Merge Engine**. Each user's UID is the primary key. `total_clicks` is summed and `unique_visitors` accumulates a [RoaringBitmap](https://roaringbitmap.org/) of all UIDs seen — both computed directly at the storage layer. + +```sql +CREATE TABLE user_profiles ( + uid INT, + total_clicks BIGINT, + unique_visitors BYTES, + PRIMARY KEY (uid) NOT ENFORCED +) WITH ( + 'table.merge-engine' = 'aggregation', + 'fields.total_clicks.agg' = 'sum', + 'fields.unique_visitors.agg' = 'rbm32' +); +``` + +## Step 4: Ingest and Process Data + +Create a temporary source table to simulate raw click events using the Faker connector. + +:::note +Java Faker's `numberBetween(min, max)` treats `max` as exclusive. The expression below produces click counts of 1–10. +::: + +```sql +CREATE TEMPORARY TABLE raw_events ( + email STRING, + click_count INT, + proctime AS PROCTIME() +) WITH ( + 'connector' = 'faker', + 'rows-per-second' = '1', + 'fields.email.expression' = '#{internet.emailAddress}', + 'fields.click_count.expression' = '#{number.numberBetween ''1'',''11''}' +); +``` + +Now run the pipeline. The `lookup.insert-if-not-exists` hint ensures that if an email is not found in `user_dict`, Fluss generates a new `uid` automatically. `rb_build_agg(d.uid)` builds a one-element RoaringBitmap from each UID — the Aggregation Merge Engine OR-s it into the stored bitmap, giving an exact unique visitor count per user over time. + +```sql +INSERT INTO user_profiles +SELECT + d.uid, + CAST(e.click_count AS BIGINT), + rb_build_agg(d.uid) +FROM raw_events AS e +JOIN user_dict /*+ OPTIONS('lookup.insert-if-not-exists' = 'true') */ +FOR SYSTEM_TIME AS OF e.proctime AS d +ON e.email = d.email +GROUP BY d.uid, e.click_count; +``` + +## Step 5: Verify Results + +Open a **second terminal**, navigate to the working directory, and launch another SQL Client session to query results while the pipeline runs. Review Comment: Why do we need a second terminal here? Flink SQL Client executes DML asynchronously by default (`table.dml-sync=false`), so the prompt becomes available again after the `INSERT INTO` job is submitted. Since this tutorial never enables synchronous DML, we can run the verification queries in the same session and remove the duplicated SQL Client startup and catalog setup. ########## website/docs/quickstart/user_profile.md: ########## @@ -0,0 +1,290 @@ +--- +title: Real-Time User Profile +sidebar_position: 4 +--- + +# Real-Time User Profile + +This tutorial demonstrates how to build a real-time user profiling system using three core Apache Fluss features: the **Auto-Increment Column**, the **Aggregation Merge Engine**, and the built-in **RoaringBitmap SQL functions**. You will learn how to automatically map high-cardinality email identifiers to compact integer UIDs, accumulate click metrics, and count unique visitors — all directly in the storage layer, keeping the Flink job entirely stateless. + +## How the System Works + +### Core Concepts + +- **Identity Mapping**: Incoming email strings are automatically mapped to compact `INT` UIDs using Fluss's auto-increment column — no manual ID management required. +- **Storage-Level Aggregation**: Click counts are summed and unique visitor bitmaps are OR-ed directly inside the Fluss TabletServers via the Aggregation Merge Engine. +- **Built-in Bitmap Functions**: `rb_build_agg` and `rb_cardinality` are registered natively in FlussCatalog — no external JAR or `CREATE TEMPORARY FUNCTION` required. + +### Data Flow + +1. **Ingestion**: Raw click events arrive with an email address and a click count. +2. **Mapping**: A Flink lookup join against `user_dict` resolves the email to a UID. If the email is new, the `insert-if-not-exists` hint instructs Fluss to generate a new UID automatically. +3. **Aggregation**: The resolved UID is written to `user_profiles`. The Aggregation Merge Engine sums `total_clicks` and OR-s the `unique_visitors` bitmap at the storage layer — no windowing or Flink state required. + +## Prerequisites + +Before proceeding, ensure that [Docker](https://docs.docker.com/engine/install/) and the [Docker Compose plugin](https://docs.docker.com/compose/install/linux/) are installed on your machine. + +## Environment Setup + +1. Create a working directory and navigate into it. + ```shell + mkdir fluss-user-profile + cd fluss-user-profile + ``` + +2. Create a `docker-compose.yml` file with the following content: + ```yaml + services: + coordinator-server: + image: apache/fluss:$FLUSS_DOCKER_VERSION$ + command: coordinatorServer + depends_on: + - zookeeper + environment: + - | + FLUSS_PROPERTIES= + zookeeper.address: zookeeper:2181 + bind.listeners: FLUSS://coordinator-server:9123 + remote.data.dir: /tmp/fluss/remote + volumes: + - fluss-remote-data:/tmp/fluss/remote + tablet-server: + image: apache/fluss:$FLUSS_DOCKER_VERSION$ + command: tabletServer + depends_on: + - coordinator-server + environment: + - | + FLUSS_PROPERTIES= + zookeeper.address: zookeeper:2181 + bind.listeners: FLUSS://tablet-server:9123 + data.dir: /tmp/fluss/data + remote.data.dir: /tmp/fluss/remote + volumes: + - fluss-remote-data:/tmp/fluss/remote + zookeeper: + restart: always + image: zookeeper:3.9.2 + jobmanager: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + ports: + - "8083:8081" + command: jobmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.port: 8081 + volumes: + - fluss-remote-data:/tmp/fluss/remote + taskmanager: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + depends_on: + - jobmanager + command: taskmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + taskmanager.numberOfTaskSlots: 2 + volumes: + - fluss-remote-data:/tmp/fluss/remote + sql-client: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + command: ["/opt/sql-client/sql-client"] + depends_on: + - jobmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.port: 8081 + volumes: + - fluss-remote-data:/tmp/fluss/remote + + volumes: + fluss-remote-data: + ``` + +3. Start all services. + ```shell + docker compose up -d + ``` + +4. Confirm all containers are running. + ```shell + docker compose ps + ``` + You should see `coordinator-server`, `tablet-server`, `zookeeper`, `jobmanager`, `taskmanager`, and `sql-client` all in the `running` state. + +:::note +All the following commands involving `docker compose` should be executed in the working directory that contains the `docker-compose.yml` file. +::: + +## Enter the SQL Client + +Use the following command to enter the Flink SQL Client: + +```shell +docker compose run --entrypoint bash sql-client -c " +\${FLINK_HOME}/bin/sql-client.sh \ + -Drest.address=jobmanager \ + -Drest.port=8081 \ + -i /opt/sql-client/sql/sql-client.sql +" +``` + +## Step 1: Create the Fluss Catalog + +Run these statements one by one in the SQL Client. + +:::tip +Run SQL statements one by one to avoid errors. +::: + +```sql +CREATE CATALOG fluss_catalog WITH ( + 'type' = 'fluss', + 'bootstrap.servers' = 'coordinator-server:9123' +); +``` + +```sql +USE CATALOG fluss_catalog; +``` + +:::note +Once you switch to the Fluss catalog, all RoaringBitmap SQL functions (`rb_build_agg`, `rb_cardinality`, `rb_or_agg`, and others) are available immediately — no `CREATE TEMPORARY FUNCTION` statement is needed. +::: + +## Step 2: Create the User Dictionary Table + +Create the `user_dict` table to map email addresses to integer UIDs. The `auto-increment.fields` property instructs Fluss to automatically assign a unique `INT` UID for every new email it receives. + +```sql +CREATE TABLE user_dict ( + email STRING, + uid INT, + PRIMARY KEY (email) NOT ENFORCED +) WITH ( + 'auto-increment.fields' = 'uid' +); +``` + +## Step 3: Create the Aggregated Profile Table + +Create the `user_profiles` table using the **Aggregation Merge Engine**. Each user's UID is the primary key. `total_clicks` is summed and `unique_visitors` accumulates a [RoaringBitmap](https://roaringbitmap.org/) of all UIDs seen — both computed directly at the storage layer. + +```sql +CREATE TABLE user_profiles ( + uid INT, + total_clicks BIGINT, + unique_visitors BYTES, + PRIMARY KEY (uid) NOT ENFORCED +) WITH ( + 'table.merge-engine' = 'aggregation', + 'fields.total_clicks.agg' = 'sum', + 'fields.unique_visitors.agg' = 'rbm32' +); +``` + +## Step 4: Ingest and Process Data + +Create a temporary source table to simulate raw click events using the Faker connector. + +:::note +Java Faker's `numberBetween(min, max)` treats `max` as exclusive. The expression below produces click counts of 1–10. +::: + +```sql +CREATE TEMPORARY TABLE raw_events ( + email STRING, + click_count INT, + proctime AS PROCTIME() +) WITH ( + 'connector' = 'faker', + 'rows-per-second' = '1', + 'fields.email.expression' = '#{internet.emailAddress}', + 'fields.click_count.expression' = '#{number.numberBetween ''1'',''11''}' +); +``` + +Now run the pipeline. The `lookup.insert-if-not-exists` hint ensures that if an email is not found in `user_dict`, Fluss generates a new `uid` automatically. `rb_build_agg(d.uid)` builds a one-element RoaringBitmap from each UID — the Aggregation Merge Engine OR-s it into the stored bitmap, giving an exact unique visitor count per user over time. + +```sql +INSERT INTO user_profiles +SELECT + d.uid, + CAST(e.click_count AS BIGINT), + rb_build_agg(d.uid) +FROM raw_events AS e +JOIN user_dict /*+ OPTIONS('lookup.insert-if-not-exists' = 'true') */ +FOR SYSTEM_TIME AS OF e.proctime AS d +ON e.email = d.email +GROUP BY d.uid, e.click_count; Review Comment: Using `rb_build_agg` here introduces an unbounded Flink group aggregation keyed by `(uid, click_count)`. Since `uid` is high-cardinality and these groups never close, Flink must retain continuously growing keyed state, which defeats the purpose of offloading aggregation to the Fluss Aggregation Merge Engine. We can remove the `GROUP BY` and use the scalar `rb_build(ARRAY[d.uid])` to emit a singleton bitmap for each event; the Fluss `rbm32` field will union these bitmaps, while `sum` aggregates the click counts on the server side. ########## website/docs/quickstart/user_profile.md: ########## @@ -0,0 +1,290 @@ +--- +title: Real-Time User Profile +sidebar_position: 4 +--- + +# Real-Time User Profile + +This tutorial demonstrates how to build a real-time user profiling system using three core Apache Fluss features: the **Auto-Increment Column**, the **Aggregation Merge Engine**, and the built-in **RoaringBitmap SQL functions**. You will learn how to automatically map high-cardinality email identifiers to compact integer UIDs, accumulate click metrics, and count unique visitors — all directly in the storage layer, keeping the Flink job entirely stateless. + +## How the System Works + +### Core Concepts + +- **Identity Mapping**: Incoming email strings are automatically mapped to compact `INT` UIDs using Fluss's auto-increment column — no manual ID management required. +- **Storage-Level Aggregation**: Click counts are summed and unique visitor bitmaps are OR-ed directly inside the Fluss TabletServers via the Aggregation Merge Engine. +- **Built-in Bitmap Functions**: `rb_build_agg` and `rb_cardinality` are registered natively in FlussCatalog — no external JAR or `CREATE TEMPORARY FUNCTION` required. + +### Data Flow + +1. **Ingestion**: Raw click events arrive with an email address and a click count. +2. **Mapping**: A Flink lookup join against `user_dict` resolves the email to a UID. If the email is new, the `insert-if-not-exists` hint instructs Fluss to generate a new UID automatically. +3. **Aggregation**: The resolved UID is written to `user_profiles`. The Aggregation Merge Engine sums `total_clicks` and OR-s the `unique_visitors` bitmap at the storage layer — no windowing or Flink state required. + +## Prerequisites + +Before proceeding, ensure that [Docker](https://docs.docker.com/engine/install/) and the [Docker Compose plugin](https://docs.docker.com/compose/install/linux/) are installed on your machine. + +## Environment Setup + +1. Create a working directory and navigate into it. + ```shell + mkdir fluss-user-profile + cd fluss-user-profile + ``` + +2. Create a `docker-compose.yml` file with the following content: + ```yaml + services: + coordinator-server: + image: apache/fluss:$FLUSS_DOCKER_VERSION$ + command: coordinatorServer + depends_on: + - zookeeper + environment: + - | + FLUSS_PROPERTIES= + zookeeper.address: zookeeper:2181 + bind.listeners: FLUSS://coordinator-server:9123 + remote.data.dir: /tmp/fluss/remote + volumes: + - fluss-remote-data:/tmp/fluss/remote + tablet-server: + image: apache/fluss:$FLUSS_DOCKER_VERSION$ + command: tabletServer + depends_on: + - coordinator-server + environment: + - | + FLUSS_PROPERTIES= + zookeeper.address: zookeeper:2181 + bind.listeners: FLUSS://tablet-server:9123 + data.dir: /tmp/fluss/data + remote.data.dir: /tmp/fluss/remote + volumes: + - fluss-remote-data:/tmp/fluss/remote + zookeeper: + restart: always + image: zookeeper:3.9.2 + jobmanager: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + ports: + - "8083:8081" + command: jobmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.port: 8081 + volumes: + - fluss-remote-data:/tmp/fluss/remote + taskmanager: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + depends_on: + - jobmanager + command: taskmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + taskmanager.numberOfTaskSlots: 2 + volumes: + - fluss-remote-data:/tmp/fluss/remote + sql-client: + image: apache/fluss-quickstart-flink:$FLUSS_QUICKSTART_FLINK_DOCKER_VERSION$ + command: ["/opt/sql-client/sql-client"] + depends_on: + - jobmanager + environment: + - | + FLINK_PROPERTIES= + jobmanager.rpc.address: jobmanager + rest.address: jobmanager + rest.port: 8081 + volumes: + - fluss-remote-data:/tmp/fluss/remote + + volumes: + fluss-remote-data: + ``` + +3. Start all services. + ```shell + docker compose up -d + ``` + +4. Confirm all containers are running. + ```shell + docker compose ps + ``` + You should see `coordinator-server`, `tablet-server`, `zookeeper`, `jobmanager`, `taskmanager`, and `sql-client` all in the `running` state. + +:::note +All the following commands involving `docker compose` should be executed in the working directory that contains the `docker-compose.yml` file. +::: + +## Enter the SQL Client + +Use the following command to enter the Flink SQL Client: + +```shell +docker compose run --entrypoint bash sql-client -c " +\${FLINK_HOME}/bin/sql-client.sh \ + -Drest.address=jobmanager \ + -Drest.port=8081 \ + -i /opt/sql-client/sql/sql-client.sql +" +``` + +## Step 1: Create the Fluss Catalog + +Run these statements one by one in the SQL Client. + +:::tip +Run SQL statements one by one to avoid errors. +::: + +```sql +CREATE CATALOG fluss_catalog WITH ( + 'type' = 'fluss', + 'bootstrap.servers' = 'coordinator-server:9123' +); +``` + +```sql +USE CATALOG fluss_catalog; +``` + +:::note +Once you switch to the Fluss catalog, all RoaringBitmap SQL functions (`rb_build_agg`, `rb_cardinality`, `rb_or_agg`, and others) are available immediately — no `CREATE TEMPORARY FUNCTION` statement is needed. +::: + +## Step 2: Create the User Dictionary Table + +Create the `user_dict` table to map email addresses to integer UIDs. The `auto-increment.fields` property instructs Fluss to automatically assign a unique `INT` UID for every new email it receives. + +```sql +CREATE TABLE user_dict ( + email STRING, + uid INT, + PRIMARY KEY (email) NOT ENFORCED +) WITH ( + 'auto-increment.fields' = 'uid' +); +``` + +## Step 3: Create the Aggregated Profile Table + +Create the `user_profiles` table using the **Aggregation Merge Engine**. Each user's UID is the primary key. `total_clicks` is summed and `unique_visitors` accumulates a [RoaringBitmap](https://roaringbitmap.org/) of all UIDs seen — both computed directly at the storage layer. + +```sql +CREATE TABLE user_profiles ( + uid INT, + total_clicks BIGINT, + unique_visitors BYTES, + PRIMARY KEY (uid) NOT ENFORCED +) WITH ( + 'table.merge-engine' = 'aggregation', + 'fields.total_clicks.agg' = 'sum', + 'fields.unique_visitors.agg' = 'rbm32' +); +``` + +## Step 4: Ingest and Process Data + +Create a temporary source table to simulate raw click events using the Faker connector. + +:::note +Java Faker's `numberBetween(min, max)` treats `max` as exclusive. The expression below produces click counts of 1–10. +::: + +```sql +CREATE TEMPORARY TABLE raw_events ( + email STRING, + click_count INT, + proctime AS PROCTIME() +) WITH ( + 'connector' = 'faker', + 'rows-per-second' = '1', + 'fields.email.expression' = '#{internet.emailAddress}', + 'fields.click_count.expression' = '#{number.numberBetween ''1'',''11''}' +); +``` + +Now run the pipeline. The `lookup.insert-if-not-exists` hint ensures that if an email is not found in `user_dict`, Fluss generates a new `uid` automatically. `rb_build_agg(d.uid)` builds a one-element RoaringBitmap from each UID — the Aggregation Merge Engine OR-s it into the stored bitmap, giving an exact unique visitor count per user over time. + +```sql +INSERT INTO user_profiles +SELECT + d.uid, + CAST(e.click_count AS BIGINT), + rb_build_agg(d.uid) Review Comment: The `user_profiles` row key and the bitmap element are both `d.uid`. For a row keyed by UID `u`, every incoming bitmap is therefore `{u}`, so the `rbm32` union remains `{u}` and `rb_cardinality(unique_visitors)` can never grow beyond 1. This does not demonstrate unique visitors as the tutorial claims. Please introduce a separate aggregation key, such as a page, campaign, or profile-group ID, as the table primary key and keep `d.uid` as the visitor ID stored in the bitmap. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
