Skip to content

Latest commit

 

History

History
2479 lines (2105 loc) · 69.3 KB

File metadata and controls

2479 lines (2105 loc) · 69.3 KB

Ticket Marketplace Analytics Engine Setup Guide

Kafka + Flink + ClickHouse + dbt + Metabase for Ticket Marketplace

Complete local development setup for a GDPR-compliant, horizontally scalable analytics engine for a ticket marketplace platform.


Architecture Overview

Ticket Events → Kafka → Flink → ClickHouse → dbt → Metabase
(Vendors/Customers)

Components:

  • Kafka: Message streaming platform for ticket events
  • Flink: Real-time stream processing for ticket transactions
  • ClickHouse: Columnar OLAP database for analytics
  • dbt: Data transformation and business metrics
  • Metabase: Data visualization for dashboards

Prerequisites

  • Docker & Docker Compose installed
  • At least 8GB RAM available for Docker
  • 20GB free disk space
  • Python 3.8+ (for dbt and event simulators)

Project Structure

ticket-analytics-engine/
├── docker-compose.yml
├── flink/
│   ├── jobs/
│   │   └── kafka_to_clickhouse.py
│   └── Dockerfile
├── clickhouse/
│   ├── init/
│   │   └── 01_init_tables.sql
│   └── config.xml
├── dbt/
│   ├── dbt_project.yml
│   ├── profiles.yml
│   ├── models/
│   │   ├── staging/
│   │   │   ├── stg_ticket_events.sql
│   │   │   ├── stg_vendor_events.sql
│   │   │   └── stg_customer_events.sql
│   │   └── marts/
│   │       ├── vendor_performance.sql
│   │       ├── ticket_sales_metrics.sql
│   │       └── customer_behavior.sql
│   └── macros/
├── simulators/
│   ├── requirements.txt
│   ├── config.yaml
│   ├── ticket_marketplace_simulator.py
│   └── README.md
├── .env
└── README.md

Step 1: Docker Compose Configuration

Create docker-compose.yml:

version: '3.8'

services:
  # Zookeeper for Kafka
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.0
    hostname: zookeeper
    container_name: zookeeper
    ports:
      - "2181:2181"
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
    volumes:
      - zookeeper-data:/var/lib/zookeeper/data
      - zookeeper-logs:/var/lib/zookeeper/log

  # Kafka Broker
  kafka:
    image: confluentinc/cp-kafka:7.5.0
    hostname: kafka
    container_name: kafka
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
      - "9093:9093"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
      KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true'
    volumes:
      - kafka-data:/var/lib/kafka/data

  # Kafka UI (Optional - for monitoring)
  kafka-ui:
    image: provectuslabs/kafka-ui:latest
    container_name: kafka-ui
    depends_on:
      - kafka
    ports:
      - "8080:8080"
    environment:
      KAFKA_CLUSTERS_0_NAME: local
      KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:29092
      KAFKA_CLUSTERS_0_ZOOKEEPER: zookeeper:2181

  # ClickHouse Server
  clickhouse:
    image: clickhouse/clickhouse-server:23.8
    hostname: clickhouse
    container_name: clickhouse
    ports:
      - "8123:8123"
      - "9000:9000"
    environment:
      CLICKHOUSE_DB: ticket_analytics
      CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT: 1
    volumes:
      - clickhouse-data:/var/lib/clickhouse
      - ./clickhouse/init:/docker-entrypoint-initdb.d
      - ./clickhouse/config.xml:/etc/clickhouse-server/config.d/custom.xml
    ulimits:
      nofile:
        soft: 262144
        hard: 262144

  # Flink JobManager
  flink-jobmanager:
    image: flink:1.18.0-scala_2.12
    hostname: flink-jobmanager
    container_name: flink-jobmanager
    ports:
      - "8081:8081"
    command: jobmanager
    environment:
      - JOB_MANAGER_RPC_ADDRESS=flink-jobmanager
    volumes:
      - ./flink/jobs:/opt/flink/jobs
      - flink-checkpoints:/tmp/flink-checkpoints

  # Flink TaskManager
  flink-taskmanager:
    image: flink:1.18.0-scala_2.12
    hostname: flink-taskmanager
    container_name: flink-taskmanager
    depends_on:
      - flink-jobmanager
    command: taskmanager
    scale: 2
    environment:
      - JOB_MANAGER_RPC_ADDRESS=flink-jobmanager
    volumes:
      - ./flink/jobs:/opt/flink/jobs
      - flink-checkpoints:/tmp/flink-checkpoints

  # Metabase
  metabase:
    image: metabase/metabase:latest
    container_name: metabase
    ports:
      - "3000:3000"
    environment:
      MB_DB_TYPE: h2
      MB_DB_FILE: /metabase-data/metabase.db
    volumes:
      - metabase-data:/metabase-data

volumes:
  zookeeper-data:
  zookeeper-logs:
  kafka-data:
  clickhouse-data:
  flink-checkpoints:
  metabase-data:

networks:
  default:
    name: ticket-analytics-network

Step 2: ClickHouse Configuration

Create clickhouse/config.xml:

<clickhouse>
    <logger>
        <level>information</level>
        <console>true</console>
    </logger>
    
    <!-- Listen on all interfaces -->
    <listen_host>0.0.0.0</listen_host>
    
    <!-- HTTP port -->
    <http_port>8123</http_port>
    
    <!-- Native protocol port -->
    <tcp_port>9000</tcp_port>
    
    <!-- Allow connections from any host -->
    <interserver_http_host>clickhouse</interserver_http_host>
</clickhouse>

Create clickhouse/init/01_init_tables.sql:

-- Create database
CREATE DATABASE IF NOT EXISTS ticket_analytics;

-- ============================================================================
-- VENDOR TABLES
-- ============================================================================

-- Vendors dimension table
CREATE TABLE IF NOT EXISTS ticket_analytics.vendors (
    vendor_id String,
    vendor_name String,
    vendor_category String,
    vendor_country String,
    vendor_city String,
    commission_rate Decimal(5,2),
    status String DEFAULT 'active',
    created_at DateTime DEFAULT now(),
    updated_at DateTime DEFAULT now()
) ENGINE = ReplacingMergeTree(updated_at)
ORDER BY vendor_id
SETTINGS index_granularity = 8192;

-- ============================================================================
-- CUSTOMER TABLES
-- ============================================================================

-- Customers dimension table
CREATE TABLE IF NOT EXISTS ticket_analytics.customers (
    customer_id String,
    customer_email String,
    customer_name String,
    customer_country String,
    customer_city String,
    registration_date DateTime,
    user_consent Boolean DEFAULT true,
    last_purchase_date Nullable(DateTime),
    total_purchases UInt32 DEFAULT 0,
    total_spent Decimal(10,2) DEFAULT 0,
    created_at DateTime DEFAULT now(),
    updated_at DateTime DEFAULT now()
) ENGINE = ReplacingMergeTree(updated_at)
ORDER BY customer_id
SETTINGS index_granularity = 8192;

-- ============================================================================
-- EVENT TABLES
-- ============================================================================

-- Events dimension table
CREATE TABLE IF NOT EXISTS ticket_analytics.events (
    event_id String,
    vendor_id String,
    event_name String,
    event_category String,
    event_venue String,
    event_city String,
    event_country String,
    event_date DateTime,
    total_tickets UInt32,
    available_tickets UInt32,
    base_price Decimal(10,2),
    created_at DateTime DEFAULT now(),
    updated_at DateTime DEFAULT now()
) ENGINE = ReplacingMergeTree(updated_at)
ORDER BY (vendor_id, event_id)
SETTINGS index_granularity = 8192;

-- ============================================================================
-- TICKET TRANSACTION TABLES
-- ============================================================================

-- Raw ticket purchase events (from Kafka)
CREATE TABLE IF NOT EXISTS ticket_analytics.ticket_purchases_raw (
    transaction_id UUID,
    customer_id String,
    vendor_id String,
    event_id String,
    ticket_type String,
    quantity UInt16,
    unit_price Decimal(10,2),
    total_price Decimal(10,2),
    currency String DEFAULT 'USD',
    payment_method String,
    purchase_timestamp DateTime64(3),
    session_id String,
    user_consent Boolean DEFAULT true,
    ip_address Nullable(String),
    user_agent Nullable(String),
    created_at DateTime DEFAULT now()
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(purchase_timestamp)
ORDER BY (vendor_id, customer_id, purchase_timestamp)
SETTINGS index_granularity = 8192;

-- Ticket views/browsing events
CREATE TABLE IF NOT EXISTS ticket_analytics.ticket_views_raw (
    view_id UUID,
    customer_id String,
    vendor_id String,
    event_id String,
    view_timestamp DateTime64(3),
    view_duration_seconds UInt32,
    session_id String,
    user_consent Boolean DEFAULT true,
    created_at DateTime DEFAULT now()
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(view_timestamp)
ORDER BY (vendor_id, event_id, view_timestamp)
SETTINGS index_granularity = 8192;

-- Add to cart events
CREATE TABLE IF NOT EXISTS ticket_analytics.cart_events_raw (
    cart_event_id UUID,
    customer_id String,
    vendor_id String,
    event_id String,
    action String, -- 'add', 'remove', 'checkout'
    quantity UInt16,
    event_timestamp DateTime64(3),
    session_id String,
    user_consent Boolean DEFAULT true,
    created_at DateTime DEFAULT now()
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(event_timestamp)
ORDER BY (customer_id, event_timestamp)
SETTINGS index_granularity = 8192;

-- Search events
CREATE TABLE IF NOT EXISTS ticket_analytics.search_events_raw (
    search_id UUID,
    customer_id String,
    search_query String,
    search_filters String, -- JSON
    results_count UInt32,
    search_timestamp DateTime64(3),
    session_id String,
    user_consent Boolean DEFAULT true,
    created_at DateTime DEFAULT now()
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(search_timestamp)
ORDER BY (search_timestamp, customer_id)
SETTINGS index_granularity = 8192;

-- ============================================================================
-- AGGREGATED TABLES (Materialized Views Targets)
-- ============================================================================

-- Vendor daily sales summary
CREATE TABLE IF NOT EXISTS ticket_analytics.vendor_daily_sales (
    vendor_id String,
    sale_date Date,
    total_transactions UInt32,
    total_tickets_sold UInt32,
    total_revenue Decimal(15,2),
    unique_customers UInt32,
    avg_transaction_value Decimal(10,2),
    created_at DateTime DEFAULT now()
) ENGINE = SummingMergeTree()
PARTITION BY toYYYYMM(sale_date)
ORDER BY (vendor_id, sale_date)
SETTINGS index_granularity = 8192;

-- Event daily sales summary
CREATE TABLE IF NOT EXISTS ticket_analytics.event_daily_sales (
    event_id String,
    vendor_id String,
    sale_date Date,
    total_transactions UInt32,
    total_tickets_sold UInt32,
    total_revenue Decimal(15,2),
    unique_customers UInt32,
    created_at DateTime DEFAULT now()
) ENGINE = SummingMergeTree()
PARTITION BY toYYYYMM(sale_date)
ORDER BY (event_id, sale_date)
SETTINGS index_granularity = 8192;

-- Customer purchase summary
CREATE TABLE IF NOT EXISTS ticket_analytics.customer_purchase_summary (
    customer_id String,
    purchase_date Date,
    total_transactions UInt32,
    total_tickets_purchased UInt32,
    total_spent Decimal(15,2),
    unique_vendors UInt32,
    unique_events UInt32,
    created_at DateTime DEFAULT now()
) ENGINE = SummingMergeTree()
PARTITION BY toYYYYMM(purchase_date)
ORDER BY (customer_id, purchase_date)
SETTINGS index_granularity = 8192;

-- Hourly marketplace metrics
CREATE TABLE IF NOT EXISTS ticket_analytics.marketplace_hourly_metrics (
    metric_hour DateTime,
    total_transactions UInt32,
    total_revenue Decimal(15,2),
    total_tickets_sold UInt32,
    unique_customers UInt32,
    unique_vendors UInt32,
    avg_transaction_value Decimal(10,2),
    created_at DateTime DEFAULT now()
) ENGINE = SummingMergeTree()
PARTITION BY toYYYYMM(metric_hour)
ORDER BY metric_hour
SETTINGS index_granularity = 8192;

-- ============================================================================
-- MATERIALIZED VIEWS (Real-time Aggregations)
-- ============================================================================

-- Vendor daily sales materialized view
CREATE MATERIALIZED VIEW IF NOT EXISTS ticket_analytics.vendor_daily_sales_mv
TO ticket_analytics.vendor_daily_sales
AS SELECT
    vendor_id,
    toDate(purchase_timestamp) as sale_date,
    count() as total_transactions,
    sum(quantity) as total_tickets_sold,
    sum(total_price) as total_revenue,
    uniqExact(customer_id) as unique_customers,
    avg(total_price) as avg_transaction_value
FROM ticket_analytics.ticket_purchases_raw
WHERE user_consent = true
GROUP BY vendor_id, sale_date;

-- Event daily sales materialized view
CREATE MATERIALIZED VIEW IF NOT EXISTS ticket_analytics.event_daily_sales_mv
TO ticket_analytics.event_daily_sales
AS SELECT
    event_id,
    vendor_id,
    toDate(purchase_timestamp) as sale_date,
    count() as total_transactions,
    sum(quantity) as total_tickets_sold,
    sum(total_price) as total_revenue,
    uniqExact(customer_id) as unique_customers
FROM ticket_analytics.ticket_purchases_raw
WHERE user_consent = true
GROUP BY event_id, vendor_id, sale_date;

-- Customer purchase summary materialized view
CREATE MATERIALIZED VIEW IF NOT EXISTS ticket_analytics.customer_purchase_summary_mv
TO ticket_analytics.customer_purchase_summary
AS SELECT
    customer_id,
    toDate(purchase_timestamp) as purchase_date,
    count() as total_transactions,
    sum(quantity) as total_tickets_purchased,
    sum(total_price) as total_spent,
    uniqExact(vendor_id) as unique_vendors,
    uniqExact(event_id) as unique_events
FROM ticket_analytics.ticket_purchases_raw
WHERE user_consent = true
GROUP BY customer_id, purchase_date;

-- Hourly marketplace metrics materialized view
CREATE MATERIALIZED VIEW IF NOT EXISTS ticket_analytics.marketplace_hourly_metrics_mv
TO ticket_analytics.marketplace_hourly_metrics
AS SELECT
    toStartOfHour(purchase_timestamp) as metric_hour,
    count() as total_transactions,
    sum(total_price) as total_revenue,
    sum(quantity) as total_tickets_sold,
    uniqExact(customer_id) as unique_customers,
    uniqExact(vendor_id) as unique_vendors,
    avg(total_price) as avg_transaction_value
FROM ticket_analytics.ticket_purchases_raw
WHERE user_consent = true
GROUP BY metric_hour;

-- ============================================================================
-- CONVERSION FUNNEL TABLE
-- ============================================================================

-- Funnel analysis: View -> Add to Cart -> Purchase
CREATE TABLE IF NOT EXISTS ticket_analytics.conversion_funnel (
    funnel_date Date,
    vendor_id String,
    event_id String,
    total_views UInt32,
    total_add_to_cart UInt32,
    total_purchases UInt32,
    view_to_cart_rate Decimal(5,2),
    cart_to_purchase_rate Decimal(5,2),
    overall_conversion_rate Decimal(5,2),
    created_at DateTime DEFAULT now()
) ENGINE = ReplacingMergeTree(created_at)
PARTITION BY toYYYYMM(funnel_date)
ORDER BY (vendor_id, event_id, funnel_date)
SETTINGS index_granularity = 8192;

-- ============================================================================
-- GDPR COMPLIANCE TABLES
-- ============================================================================

-- User data deletion requests
CREATE TABLE IF NOT EXISTS ticket_analytics.user_deletions (
    customer_id String,
    deletion_requested_at DateTime,
    deletion_completed_at Nullable(DateTime),
    deletion_status String DEFAULT 'pending',
    created_at DateTime DEFAULT now()
) ENGINE = ReplacingMergeTree(created_at)
ORDER BY (customer_id, deletion_requested_at)
SETTINGS index_granularity = 8192;

-- Data access audit log
CREATE TABLE IF NOT EXISTS ticket_analytics.data_access_log (
    log_id UUID,
    customer_id String,
    access_type String, -- 'export', 'view', 'delete'
    accessed_by String,
    access_timestamp DateTime,
    ip_address String,
    created_at DateTime DEFAULT now()
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(access_timestamp)
ORDER BY (customer_id, access_timestamp)
SETTINGS index_granularity = 8192;

Step 3: Flink Job Configuration

Create flink/jobs/kafka_to_clickhouse.py (PyFlink example):

from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment, EnvironmentSettings

def create_kafka_purchase_source(t_env):
    t_env.execute_sql("""
        CREATE TABLE kafka_ticket_purchases (
            transaction_id STRING,
            customer_id STRING,
            vendor_id STRING,
            event_id STRING,
            ticket_type STRING,
            quantity INT,
            unit_price DECIMAL(10,2),
            total_price DECIMAL(10,2),
            currency STRING,
            payment_method STRING,
            purchase_timestamp TIMESTAMP(3),
            session_id STRING,
            user_consent BOOLEAN,
            ip_address STRING,
            user_agent STRING,
            WATERMARK FOR purchase_timestamp AS purchase_timestamp - INTERVAL '5' SECOND
        ) WITH (
            'connector' = 'kafka',
            'topic' = 'ticket-purchases',
            'properties.bootstrap.servers' = 'kafka:29092',
            'properties.group.id' = 'flink-purchases-consumer',
            'scan.startup.mode' = 'latest-offset',
            'format' = 'json',
            'json.fail-on-missing-field' = 'false',
            'json.ignore-parse-errors' = 'true'
        )
    """)

def create_kafka_view_source(t_env):
    t_env.execute_sql("""
        CREATE TABLE kafka_ticket_views (
            view_id STRING,
            customer_id STRING,
            vendor_id STRING,
            event_id STRING,
            view_timestamp TIMESTAMP(3),
            view_duration_seconds INT,
            session_id STRING,
            user_consent BOOLEAN,
            WATERMARK FOR view_timestamp AS view_timestamp - INTERVAL '5' SECOND
        ) WITH (
            'connector' = 'kafka',
            'topic' = 'ticket-views',
            'properties.bootstrap.servers' = 'kafka:29092',
            'properties.group.id' = 'flink-views-consumer',
            'scan.startup.mode' = 'latest-offset',
            'format' = 'json',
            'json.fail-on-missing-field' = 'false',
            'json.ignore-parse-errors' = 'true'
        )
    """)

def create_kafka_cart_source(t_env):
    t_env.execute_sql("""
        CREATE TABLE kafka_cart_events (
            cart_event_id STRING,
            customer_id STRING,
            vendor_id STRING,
            event_id STRING,
            action STRING,
            quantity INT,
            event_timestamp TIMESTAMP(3),
            session_id STRING,
            user_consent BOOLEAN,
            WATERMARK FOR event_timestamp AS event_timestamp - INTERVAL '5' SECOND
        ) WITH (
            'connector' = 'kafka',
            'topic' = 'cart-events',
            'properties.bootstrap.servers' = 'kafka:29092',
            'properties.group.id' = 'flink-cart-consumer',
            'scan.startup.mode' = 'latest-offset',
            'format' = 'json',
            'json.fail-on-missing-field' = 'false',
            'json.ignore-parse-errors' = 'true'
        )
    """)

def create_clickhouse_purchase_sink(t_env):
    t_env.execute_sql("""
        CREATE TABLE clickhouse_ticket_purchases (
            transaction_id STRING,
            customer_id STRING,
            vendor_id STRING,
            event_id STRING,
            ticket_type STRING,
            quantity INT,
            unit_price DECIMAL(10,2),
            total_price DECIMAL(10,2),
            currency STRING,
            payment_method STRING,
            purchase_timestamp TIMESTAMP(3),
            session_id STRING,
            user_consent BOOLEAN,
            ip_address STRING,
            user_agent STRING
        ) WITH (
            'connector' = 'jdbc',
            'url' = 'jdbc:clickhouse://clickhouse:8123/ticket_analytics',
            'table-name' = 'ticket_purchases_raw',
            'driver' = 'com.clickhouse.jdbc.ClickHouseDriver',
            'sink.buffer-flush.max-rows' = '1000',
            'sink.buffer-flush.interval' = '5s'
        )
    """)

def create_clickhouse_view_sink(t_env):
    t_env.execute_sql("""
        CREATE TABLE clickhouse_ticket_views (
            view_id STRING,
            customer_id STRING,
            vendor_id STRING,
            event_id STRING,
            view_timestamp TIMESTAMP(3),
            view_duration_seconds INT,
            session_id STRING,
            user_consent BOOLEAN
        ) WITH (
            'connector' = 'jdbc',
            'url' = 'jdbc:clickhouse://clickhouse:8123/ticket_analytics',
            'table-name' = 'ticket_views_raw',
            'driver' = 'com.clickhouse.jdbc.ClickHouseDriver',
            'sink.buffer-flush.max-rows' = '1000',
            'sink.buffer-flush.interval' = '5s'
        )
    """)

def create_clickhouse_cart_sink(t_env):
    t_env.execute_sql("""
        CREATE TABLE clickhouse_cart_events (
            cart_event_id STRING,
            customer_id STRING,
            vendor_id STRING,
            event_id STRING,
            action STRING,
            quantity INT,
            event_timestamp TIMESTAMP(3),
            session_id STRING,
            user_consent BOOLEAN
        ) WITH (
            'connector' = 'jdbc',
            'url' = 'jdbc:clickhouse://clickhouse:8123/ticket_analytics',
            'table-name' = 'cart_events_raw',
            'driver' = 'com.clickhouse.jdbc.ClickHouseDriver',
            'sink.buffer-flush.max-rows' = '1000',
            'sink.buffer-flush.interval' = '5s'
        )
    """)

def main():
    # Create execution environment
    env = StreamExecutionEnvironment.get_execution_environment()
    env.set_parallelism(2)
    
    # Create table environment
    settings = EnvironmentSettings.new_instance().in_streaming_mode().build()
    t_env = StreamTableEnvironment.create(env, environment_settings=settings)
    
    # Create source and sink tables
    create_kafka_purchase_source(t_env)
    create_kafka_view_source(t_env)
    create_kafka_cart_source(t_env)
    
    create_clickhouse_purchase_sink(t_env)
    create_clickhouse_view_sink(t_env)
    create_clickhouse_cart_sink(t_env)
    
    # GDPR filter: Only process events with user consent
    # Purchases
    t_env.execute_sql("""
        INSERT INTO clickhouse_ticket_purchases
        SELECT 
            transaction_id,
            customer_id,
            vendor_id,
            event_id,
            ticket_type,
            quantity,
            unit_price,
            total_price,
            currency,
            payment_method,
            purchase_timestamp,
            session_id,
            user_consent,
            ip_address,
            user_agent
        FROM kafka_ticket_purchases
        WHERE user_consent = true
    """)
    
    # Views
    t_env.execute_sql("""
        INSERT INTO clickhouse_ticket_views
        SELECT 
            view_id,
            customer_id,
            vendor_id,
            event_id,
            view_timestamp,
            view_duration_seconds,
            session_id,
            user_consent
        FROM kafka_ticket_views
        WHERE user_consent = true
    """)
    
    # Cart events
    t_env.execute_sql("""
        INSERT INTO clickhouse_cart_events
        SELECT 
            cart_event_id,
            customer_id,
            vendor_id,
            event_id,
            action,
            quantity,
            event_timestamp,
            session_id,
            user_consent
        FROM kafka_cart_events
        WHERE user_consent = true
    """)

if __name__ == '__main__':
    main()

Step 4: dbt Configuration

Create dbt/dbt_project.yml:

name: 'ticket_analytics'
version: '1.0.0'
config-version: 2

profile: 'ticket_analytics'

model-paths: ["models"]
analysis-paths: ["analyses"]
test-paths: ["tests"]
seed-paths: ["seeds"]
macro-paths: ["macros"]
snapshot-paths: ["snapshots"]

target-path: "target"
clean-targets:
  - "target"
  - "dbt_packages"

models:
  ticket_analytics:
    staging:
      +materialized: view
    marts:
      +materialized: table

Create dbt/profiles.yml:

ticket_analytics:
  target: dev
  outputs:
    dev:
      type: clickhouse
      host: localhost
      port: 8123
      user: default
      password: ''
      database: ticket_analytics
      schema: ticket_analytics
      secure: false
      verify: false

Create dbt/models/staging/stg_ticket_purchases.sql:

-- Staging model: Clean and standardize ticket purchases
WITH source AS (
    SELECT * FROM {{ source('ticket_analytics', 'ticket_purchases_raw') }}
    WHERE user_consent = true
      AND purchase_timestamp >= now() - INTERVAL 90 DAY  -- GDPR: 90-day retention
),

cleaned AS (
    SELECT
        transaction_id,
        customer_id,
        vendor_id,
        event_id,
        ticket_type,
        quantity,
        unit_price,
        total_price,
        currency,
        payment_method,
        purchase_timestamp,
        session_id,
        toDate(purchase_timestamp) AS purchase_date,
        toHour(purchase_timestamp) AS purchase_hour,
        toDayOfWeek(purchase_timestamp) AS day_of_week,
        created_at
    FROM source
)

SELECT * FROM cleaned

Create dbt/models/staging/stg_ticket_views.sql:

-- Staging model: Clean and standardize ticket views
WITH source AS (
    SELECT * FROM {{ source('ticket_analytics', 'ticket_views_raw') }}
    WHERE user_consent = true
      AND view_timestamp >= now() - INTERVAL 90 DAY
),

cleaned AS (
    SELECT
        view_id,
        customer_id,
        vendor_id,
        event_id,
        view_timestamp,
        view_duration_seconds,
        session_id,
        toDate(view_timestamp) AS view_date,
        created_at
    FROM source
)

SELECT * FROM cleaned

Create dbt/models/marts/vendor_performance.sql:

-- Mart model: Vendor performance metrics
{{ config(materialized='table') }}

WITH vendor_sales AS (
    SELECT
        v.vendor_id,
        v.vendor_name,
        v.vendor_category,
        v.vendor_country,
        COUNT(DISTINCT p.transaction_id) AS total_transactions,
        SUM(p.quantity) AS total_tickets_sold,
        SUM(p.total_price) AS total_revenue,
        AVG(p.total_price) AS avg_transaction_value,
        COUNT(DISTINCT p.customer_id) AS unique_customers,
        COUNT(DISTINCT p.event_id) AS events_with_sales,
        MIN(p.purchase_timestamp) AS first_sale_date,
        MAX(p.purchase_timestamp) AS last_sale_date
    FROM {{ source('ticket_analytics', 'vendors') }} v
    LEFT JOIN {{ ref('stg_ticket_purchases') }} p ON v.vendor_id = p.vendor_id
    WHERE v.status = 'active'
    GROUP BY v.vendor_id, v.vendor_name, v.vendor_category, v.vendor_country
),

vendor_views AS (
    SELECT
        vendor_id,
        COUNT(DISTINCT view_id) AS total_views,
        AVG(view_duration_seconds) AS avg_view_duration
    FROM {{ ref('stg_ticket_views') }}
    GROUP BY vendor_id
)

SELECT
    vs.vendor_id,
    vs.vendor_name,
    vs.vendor_category,
    vs.vendor_country,
    vs.total_transactions,
    vs.total_tickets_sold,
    vs.total_revenue,
    vs.avg_transaction_value,
    vs.unique_customers,
    vs.events_with_sales,
    vv.total_views,
    vv.avg_view_duration,
    CASE 
        WHEN vv.total_views > 0 THEN (vs.total_transactions::Float / vv.total_views::Float) * 100
        ELSE 0
    END AS conversion_rate,
    vs.first_sale_date,
    vs.last_sale_date,
    now() AS created_at
FROM vendor_sales vs
LEFT JOIN vendor_views vv ON vs.vendor_id = vv.vendor_id

Create dbt/models/marts/ticket_sales_metrics.sql:

-- Mart model: Ticket sales metrics by event
{{ config(materialized='table') }}

WITH event_sales AS (
    SELECT
        e.event_id,
        e.event_name,
        e.event_category,
        e.event_venue,
        e.event_city,
        e.event_date,
        e.vendor_id,
        v.vendor_name,
        e.total_tickets,
        e.base_price,
        COUNT(DISTINCT p.transaction_id) AS transactions_count,
        SUM(p.quantity) AS tickets_sold,
        SUM(p.total_price) AS revenue,
        AVG(p.unit_price) AS avg_ticket_price,
        COUNT(DISTINCT p.customer_id) AS unique_buyers,
        MIN(p.purchase_timestamp) AS first_purchase,
        MAX(p.purchase_timestamp) AS last_purchase
    FROM {{ source('ticket_analytics', 'events') }} e
    LEFT JOIN {{ source('ticket_analytics', 'vendors') }} v ON e.vendor_id = v.vendor_id
    LEFT JOIN {{ ref('stg_ticket_purchases') }} p ON e.event_id = p.event_id
    GROUP BY 
        e.event_id, e.event_name, e.event_category, e.event_venue, 
        e.event_city, e.event_date, e.vendor_id, v.vendor_name,
        e.total_tickets, e.base_price
)

SELECT
    event_id,
    event_name,
    event_category,
    event_venue,
    event_city,
    event_date,
    vendor_id,
    vendor_name,
    total_tickets,
    base_price,
    transactions_count,
    tickets_sold,
    revenue,
    avg_ticket_price,
    unique_buyers,
    CASE 
        WHEN total_tickets > 0 THEN (tickets_sold::Float / total_tickets::Float) * 100
        ELSE 0
    END AS sell_through_rate,
    CASE
        WHEN tickets_sold > 0 THEN revenue::Float / tickets_sold::Float
        ELSE 0
    END AS revenue_per_ticket,
    first_purchase,
    last_purchase,
    dateDiff('day', first_purchase, last_purchase) AS sales_duration_days,
    now() AS created_at
FROM event_sales

Create dbt/models/marts/customer_behavior.sql:

-- Mart model: Customer behavior and segmentation
{{ config(materialized='table') }}

WITH customer_purchases AS (
    SELECT
        c.customer_id,
        c.customer_name,
        c.customer_country,
        c.customer_city,
        c.registration_date,
        COUNT(DISTINCT p.transaction_id) AS total_purchases,
        SUM(p.quantity) AS total_tickets_bought,
        SUM(p.total_price) AS total_spent,
        AVG(p.total_price) AS avg_transaction_value,
        COUNT(DISTINCT p.vendor_id) AS vendors_purchased_from,
        COUNT(DISTINCT p.event_id) AS unique_events_attended,
        MIN(p.purchase_timestamp) AS first_purchase_date,
        MAX(p.purchase_timestamp) AS last_purchase_date,
        dateDiff('day', MIN(p.purchase_timestamp), MAX(p.purchase_timestamp)) AS customer_lifetime_days
    FROM {{ source('ticket_analytics', 'customers') }} c
    LEFT JOIN {{ ref('stg_ticket_purchases') }} p ON c.customer_id = p.customer_id
    WHERE c.user_consent = true
    GROUP BY 
        c.customer_id, c.customer_name, c.customer_country, 
        c.customer_city, c.registration_date
),

customer_views AS (
    SELECT
        customer_id,
        COUNT(DISTINCT view_id) AS total_views,
        AVG(view_duration_seconds) AS avg_view_duration,
        COUNT(DISTINCT event_id) AS events_viewed
    FROM {{ ref('stg_ticket_views') }}
    GROUP BY customer_id
)

SELECT
    cp.customer_id,
    cp.customer_name,
    cp.customer_country,
    cp.customer_city,
    cp.registration_date,
    cp.total_purchases,
    cp.total_tickets_bought,
    cp.total_spent,
    cp.avg_transaction_value,
    cp.vendors_purchased_from,
    cp.unique_events_attended,
    cv.total_views,
    cv.avg_view_duration,
    cv.events_viewed,
    CASE
        WHEN cv.total_views > 0 THEN (cp.total_purchases::Float / cv.total_views::Float) * 100
        ELSE 0
    END AS view_to_purchase_rate,
    cp.first_purchase_date,
    cp.last_purchase_date,
    cp.customer_lifetime_days,
    dateDiff('day', cp.last_purchase_date, now()) AS days_since_last_purchase,
    CASE
        WHEN cp.total_spent >= 1000 THEN 'VIP'
        WHEN cp.total_spent >= 500 THEN 'High Value'
        WHEN cp.total_spent >= 100 THEN 'Medium Value'
        ELSE 'Low Value'
    END AS customer_segment,
    CASE
        WHEN dateDiff('day', cp.last_purchase_date, now()) <= 30 THEN 'Active'
        WHEN dateDiff('day', cp.last_purchase_date, now()) <= 90 THEN 'At Risk'
        ELSE 'Churned'
    END AS activity_status,
    now() AS created_at
FROM customer_purchases cp
LEFT JOIN customer_views cv ON cp.customer_id = cv.customer_id

Create dbt/models/schema.yml:

version: 2

sources:
  - name: ticket_analytics
    database: ticket_analytics
    tables:
      - name: ticket_purchases_raw
      - name: ticket_views_raw
      - name: cart_events_raw
      - name: vendors
      - name: customers
      - name: events

models:
  - name: stg_ticket_purchases
    description: "Cleaned and standardized ticket purchase events"
    columns:
      - name: transaction_id
        description: "Unique transaction identifier"
        tests:
          - unique
          - not_null
      - name: customer_id
        description: "Customer identifier"
        tests:
          - not_null

  - name: stg_ticket_views
    description: "Cleaned and standardized ticket view events"
    columns:
      - name: view_id
        description: "Unique view identifier"
        tests:
          - unique
          - not_null

  - name: vendor_performance
    description: "Vendor performance metrics and KPIs"
    columns:
      - name: vendor_id
        description: "Vendor identifier"
        tests:
          - unique
          - not_null
      - name: total_revenue
        description: "Total revenue generated by vendor"

  - name: ticket_sales_metrics
    description: "Event-level ticket sales metrics"
    columns:
      - name: event_id
        description: "Event identifier"

  - name: customer_behavior
    description: "Customer behavior analysis and segmentation"
    columns:
      - name: customer_id
        description: "Customer identifier"
        tests:
          - unique
          - not_null

Step 5: Ticket Marketplace Simulator

Create simulators/requirements.txt:

kafka-python==2.0.2
faker==19.6.2
pyyaml==6.0.1

Create simulators/config.yaml:

# Ticket Marketplace Simulation Configuration

simulation:
  # Duration in minutes
  duration_minutes: 10
  
  # Vendors configuration
  vendors:
    count: 50  # Number of vendors
    categories:
      - "Music Concerts"
      - "Sports"
      - "Theater & Arts"
      - "Comedy Shows"
      - "Festivals"
      - "Family Events"
    countries:
      - "USA"
      - "UK"
      - "Canada"
      - "Australia"
      - "Germany"
  
  # Customers configuration
  customers:
    count: 1000  # Number of unique customers
    consent_rate: 0.95  # 95% of customers give consent
    countries:
      - "USA"
      - "UK"
      - "Canada"
      - "Australia"
      - "Germany"
  
  # Events configuration
  events:
    count_per_vendor: 5  # Average events per vendor
    ticket_inventory_range: [100, 5000]  # Min and max tickets per event
    base_price_range: [25, 500]  # Min and max base price
  
  # Traffic patterns
  traffic:
    # Purchase rate per minute
    purchases_per_minute: 100
    
    # View rate per minute (should be higher than purchases)
    views_per_minute: 500
    
    # Cart events per minute
    cart_events_per_minute: 200
    
    # Search events per minute
    searches_per_minute: 150
    
    # Peak hours multiplier (simulates rush times)
    peak_hours: [18, 19, 20, 21]  # 6 PM - 9 PM
    peak_multiplier: 3.0  # 3x traffic during peak hours
    
    # Conversion funnel rates
    view_to_cart_rate: 0.15  # 15% of views add to cart
    cart_to_purchase_rate: 0.60  # 60% of cart additions result in purchase
  
  # Ticket types distribution
  ticket_types:
    - name: "General Admission"
      weight: 0.50
    - name: "VIP"
      weight: 0.20
    - name: "Premium"
      weight: 0.15
    - name: "Early Bird"
      weight: 0.10
    - name: "Group Package"
      weight: 0.05
  
  # Payment methods distribution
  payment_methods:
    - name: "Credit Card"
      weight: 0.60
    - name: "PayPal"
      weight: 0.20
    - name: "Debit Card"
      weight: 0.15
    - name: "Apple Pay"
      weight: 0.05

# Kafka configuration
kafka:
  bootstrap_servers: "localhost:9092"
  topics:
    purchases: "ticket-purchases"
    views: "ticket-views"
    cart_events: "cart-events"
    searches: "search-events"

# Scenarios - Pre-defined simulation scenarios
scenarios:
  normal_day:
    purchases_per_minute: 50
    views_per_minute: 300
    duration_minutes: 60
    
  flash_sale:
    purchases_per_minute: 500
    views_per_minute: 2000
    duration_minutes: 15
    peak_multiplier: 5.0
    
  concert_announcement:
    purchases_per_minute: 200
    views_per_minute: 1500
    duration_minutes: 30
    
  weekend_rush:
    purchases_per_minute: 300
    views_per_minute: 1200
    duration_minutes: 120
    peak_multiplier: 2.5

Create simulators/ticket_marketplace_simulator.py:

import json
import time
import uuid
import random
import yaml
from datetime import datetime, timedelta
from kafka import KafkaProducer
from faker import Faker
from typing import List, Dict
import threading

fake = Faker()

class TicketMarketplaceSimulator:
    def __init__(self, config_path='config.yaml'):
        """Initialize the ticket marketplace simulator"""
        with open(config_path, 'r') as f:
            self.config = yaml.safe_load(f)
        
        # Initialize Kafka producer
        self.producer = KafkaProducer(
            bootstrap_servers=self.config['kafka']['bootstrap_servers'],
            value_serializer=lambda v: json.dumps(v).encode('utf-8')
        )
        
        # Generate static data
        self.vendors = self._generate_vendors()
        self.customers = self._generate_customers()
        self.events = self._generate_events()
        self.sessions = {}  # Track customer sessions
        
        # Statistics
        self.stats = {
            'purchases': 0,
            'views': 0,
            'cart_events': 0,
            'searches': 0
        }
        
        print(f"✓ Generated {len(self.vendors)} vendors")
        print(f"✓ Generated {len(self.customers)} customers")
        print(f"✓ Generated {len(self.events)} events")
    
    def _generate_vendors(self) -> List[Dict]:
        """Generate vendor data"""
        vendors = []
        vendor_config = self.config['simulation']['vendors']
        
        for i in range(vendor_config['count']):
            vendor = {
                'vendor_id': f'vendor_{i+1:04d}',
                'vendor_name': fake.company(),
                'vendor_category': random.choice(vendor_config['categories']),
                'vendor_country': random.choice(vendor_config['countries']),
                'vendor_city': fake.city(),
                'commission_rate': round(random.uniform(5, 15), 2),
                'status': 'active'
            }
            vendors.append(vendor)
        
        return vendors
    
    def _generate_customers(self) -> List[Dict]:
        """Generate customer data"""
        customers = []
        customer_config = self.config['simulation']['customers']
        
        for i in range(customer_config['count']):
            has_consent = random.random() < customer_config['consent_rate']
            customer = {
                'customer_id': f'customer_{i+1:05d}',
                'customer_email': fake.email(),
                'customer_name': fake.name(),
                'customer_country': random.choice(customer_config['countries']),
                'customer_city': fake.city(),
                'registration_date': (datetime.now() - timedelta(days=random.randint(1, 365))).isoformat(),
                'user_consent': has_consent
            }
            customers.append(customer)
        
        return customers
    
    def _generate_events(self) -> List[Dict]:
        """Generate event data"""
        events = []
        events_config = self.config['simulation']['events']
        
        event_id = 1
        for vendor in self.vendors:
            num_events = random.randint(
                max(1, events_config['count_per_vendor'] - 2),
                events_config['count_per_vendor'] + 2
            )
            
            for _ in range(num_events):
                event_date = datetime.now() + timedelta(days=random.randint(7, 180))
                total_tickets = random.randint(*events_config['ticket_inventory_range'])
                
                event = {
                    'event_id': f'event_{event_id:05d}',
                    'vendor_id': vendor['vendor_id'],
                    'event_name': self._generate_event_name(vendor['vendor_category']),
                    'event_category': vendor['vendor_category'],
                    'event_venue': fake.company() + ' Arena',
                    'event_city': vendor['vendor_city'],
                    'event_country': vendor['vendor_country'],
                    'event_date': event_date.isoformat(),
                    'total_tickets': total_tickets,
                    'available_tickets': total_tickets,
                    'base_price': round(random.uniform(*events_config['base_price_range']), 2)
                }
                events.append(event)
                event_id += 1
        
        return events
    
    def _generate_event_name(self, category: str) -> str:
        """Generate event name based on category"""
        names = {
            "Music Concerts": ["Rock Festival", "Jazz Night", "Pop Concert", "Classical Evening"],
            "Sports": ["Championship Game", "Derby Match", "Boxing Night", "Tennis Open"],
            "Theater & Arts": ["Broadway Show", "Opera Night", "Ballet Performance", "Art Exhibition"],
            "Comedy Shows": ["Stand-up Night", "Comedy Festival", "Improv Show", "Sketch Comedy"],
            "Festivals": ["Food Festival", "Beer Fest", "Cultural Celebration", "Summer Festival"],
            "Family Events": ["Kids Show", "Family Fun Day", "Circus", "Magic Show"]
        }
        
        event_names = names.get(category, ["Special Event"])
        return fake.name() + " - " + random.choice(event_names)
    
    def _get_session_id(self, customer_id: str) -> str:
        """Get or create session ID for customer"""
        if customer_id not in self.sessions:
            self.sessions[customer_id] = f'session_{uuid.uuid4()}'
        return self.sessions[customer_id]
    
    def _is_peak_hour(self) -> bool:
        """Check if current hour is a peak hour"""
        current_hour = datetime.now().hour
        return current_hour in self.config['simulation']['traffic']['peak_hours']
    
    def _get_rate_multiplier(self) -> float:
        """Get traffic multiplier based on peak hours"""
        if self._is_peak_hour():
            return self.config['simulation']['traffic']['peak_multiplier']
        return 1.0
    
    def _weighted_choice(self, items: List[Dict], weight_key: str = 'weight'):
        """Make a weighted random choice"""
        weights = [item[weight_key] for item in items]
        return random.choices(items, weights=weights, k=1)[0]
    
    def generate_ticket_purchase(self) -> Dict:
        """Generate a ticket purchase event"""
        customer = random.choice(self.customers)
        event = random.choice(self.events)
        vendor = next(v for v in self.vendors if v['vendor_id'] == event['vendor_id'])
        
        ticket_type = self._weighted_choice(self.config['simulation']['ticket_types'])
        payment_method = self._weighted_choice(self.config['simulation']['payment_methods'])
        
        quantity = random.randint(1, 6)
        price_multiplier = {
            "General Admission": 1.0,
            "VIP": 2.5,
            "Premium": 1.8,
            "Early Bird": 0.8,
            "Group Package": 0.9
        }
        
        unit_price = round(event['base_price'] * price_multiplier.get(ticket_type['name'], 1.0), 2)
        total_price = round(unit_price * quantity, 2)
        
        purchase = {
            'transaction_id': str(uuid.uuid4()),
            'customer_id': customer['customer_id'],
            'vendor_id': vendor['vendor_id'],
            'event_id': event['event_id'],
            'ticket_type': ticket_type['name'],
            'quantity': quantity,
            'unit_price': unit_price,
            'total_price': total_price,
            'currency': 'USD',
            'payment_method': payment_method['name'],
            'purchase_timestamp': datetime.now().isoformat(),
            'session_id': self._get_session_id(customer['customer_id']),
            'user_consent': customer['user_consent'],
            'ip_address': fake.ipv4(),
            'user_agent': fake.user_agent()
        }
        
        return purchase
    
    def generate_ticket_view(self) -> Dict:
        """Generate a ticket view event"""
        customer = random.choice(self.customers)
        event = random.choice(self.events)
        
        view = {
            'view_id': str(uuid.uuid4()),
            'customer_id': customer['customer_id'],
            'vendor_id': event['vendor_id'],
            'event_id': event['event_id'],
            'view_timestamp': datetime.now().isoformat(),
            'view_duration_seconds': random.randint(5, 300),
            'session_id': self._get_session_id(customer['customer_id']),
            'user_consent': customer['user_consent']
        }
        
        return view
    
    def generate_cart_event(self) -> Dict:
        """Generate a cart event"""
        customer = random.choice(self.customers)
        event = random.choice(self.events)
        action = random.choices(['add', 'remove', 'checkout'], weights=[0.6, 0.2, 0.2])[0]
        
        cart_event = {
            'cart_event_id': str(uuid.uuid4()),
            'customer_id': customer['customer_id'],
            'vendor_id': event['vendor_id'],
            'event_id': event['event_id'],
            'action': action,
            'quantity': random.randint(1, 6),
            'event_timestamp': datetime.now().isoformat(),
            'session_id': self._get_session_id(customer['customer_id']),
            'user_consent': customer['user_consent']
        }
        
        return cart_event
    
    def generate_search_event(self) -> Dict:
        """Generate a search event"""
        customer = random.choice(self.customers)
        
        search_terms = [
            "concert", "festival", "sports", "theater", "comedy",
            "music", "family", "weekend", "tonight", "next month"
        ]
        
        filters = {
            'category': random.choice([cat for cat in self.config['simulation']['vendors']['categories']] + [None]),
            'city': random.choice([fake.city(), None]),
            'price_range': random.choice(['0-50', '50-100', '100-200', '200+', None]),
            'date_range': random.choice(['this_week', 'this_month', 'next_month', None])
        }
        
        search = {
            'search_id': str(uuid.uuid4()),
            'customer_id': customer['customer_id'],
            'search_query': ' '.join(random.sample(search_terms, random.randint(1, 3))),
            'search_filters': json.dumps({k: v for k, v in filters.items() if v is not None}),
            'results_count': random.randint(0, 100),
            'search_timestamp': datetime.now().isoformat(),
            'session_id': self._get_session_id(customer['customer_id']),
            'user_consent': customer['user_consent']
        }
        
        return search
    
    def send_event(self, topic: str, event: Dict, event_type: str):
        """Send event to Kafka"""
        try:
            self.producer.send(topic, value=event)
            self.stats[event_type] += 1
        except Exception as e:
            print(f"Error sending {event_type}: {e}")
    
    def event_generator(self, event_func, topic: str, rate_per_minute: int, event_type: str, duration_minutes: int):
        """Generate events at specified rate"""
        multiplier = self._get_rate_multiplier()
        actual_rate = int(rate_per_minute * multiplier)
        interval = 60.0 / actual_rate if actual_rate > 0 else 1.0
        
        end_time = time.time() + (duration_minutes * 60)
        
        while time.time() < end_time:
            event = event_func()
            self.send_event(topic, event, event_type)
            time.sleep(interval)
    
    def run_simulation(self, scenario: str = None):
        """Run the marketplace simulation"""
        # Load scenario or use default config
        if scenario and scenario in self.config.get('scenarios', {}):
            sim_config = self.config['scenarios'][scenario]
            print(f"\n🎯 Running scenario: {scenario}")
        else:
            sim_config = self.config['simulation']
            print(f"\n🎯 Running default simulation")
        
        traffic = sim_config.get('traffic', self.config['simulation']['traffic'])
        duration = sim_config.get('duration_minutes', self.config['simulation']['duration_minutes'])
        
        print(f"⏱️  Duration: {duration} minutes")
        print(f"🎫 Purchase rate: {traffic['purchases_per_minute']}/min")
        print(f"👀 View rate: {traffic['views_per_minute']}/min")
        print(f"🛒 Cart event rate: {traffic['cart_events_per_minute']}/min")
        print(f"🔍 Search rate: {traffic['searches_per_minute']}/min")
        print(f"\n🚀 Starting simulation...\n")
        
        # Create threads for each event type
        threads = [
            threading.Thread(
                target=self.event_generator,
                args=(
                    self.generate_ticket_purchase,
                    self.config['kafka']['topics']['purchases'],
                    traffic['purchases_per_minute'],
                    'purchases',
                    duration
                )
            ),
            threading.Thread(
                target=self.event_generator,
                args=(
                    self.generate_ticket_view,
                    self.config['kafka']['topics']['views'],
                    traffic['views_per_minute'],
                    'views',
                    duration
                )
            ),
            threading.Thread(
                target=self.event_generator,
                args=(
                    self.generate_cart_event,
                    self.config['kafka']['topics']['cart_events'],
                    traffic['cart_events_per_minute'],
                    'cart_events',
                    duration
                )
            ),
            threading.Thread(
                target=self.event_generator,
                args=(
                    self.generate_search_event,
                    self.config['kafka']['topics']['searches'],
                    traffic['searches_per_minute'],
                    'searches',
                    duration
                )
            )
        ]
        
        # Start all threads
        for thread in threads:
            thread.daemon = True
            thread.start()
        
        # Monitor progress
        start_time = time.time()
        try:
            while any(thread.is_alive() for thread in threads):
                elapsed = (time.time() - start_time) / 60
                print(f"\r⏳ Elapsed: {elapsed:.1f}min | "
                      f"Purchases: {self.stats['purchases']} | "
                      f"Views: {self.stats['views']} | "
                      f"Cart: {self.stats['cart_events']} | "
                      f"Searches: {self.stats['searches']}", end='')
                time.sleep(1)
        except KeyboardInterrupt:
            print("\n\n⚠️  Simulation interrupted by user")
        
        # Wait for all threads to complete
        for thread in threads:
            thread.join()
        
        print(f"\n\n✅ Simulation complete!")
        print(f"\n📊 Final Statistics:")
        print(f"   Purchases: {self.stats['purchases']}")
        print(f"   Views: {self.stats['views']}")
        print(f"   Cart Events: {self.stats['cart_events']}")
        print(f"   Searches: {self.stats['searches']}")
        print(f"   Total Events: {sum(self.stats.values())}")
        
        # Close producer
        self.producer.close()


def main():
    import argparse
    
    parser = argparse.ArgumentParser(description='Ticket Marketplace Event Simulator')
    parser.add_argument('--config', default='config.yaml', help='Configuration file path')
    parser.add_argument('--scenario', help='Scenario name to run (normal_day, flash_sale, concert_announcement, weekend_rush)')
    parser.add_argument('--vendors', type=int, help='Override number of vendors')
    parser.add_argument('--customers', type=int, help='Override number of customers')
    parser.add_argument('--duration', type=int, help='Override duration in minutes')
    parser.add_argument('--purchases-per-min', type=int, help='Override purchases per minute')
    
    args = parser.parse_args()
    
    # Initialize simulator
    simulator = TicketMarketplaceSimulator(args.config)
    
    # Override configuration if specified
    if args.vendors:
        simulator.config['simulation']['vendors']['count'] = args.vendors
        simulator.vendors = simulator._generate_vendors()
        simulator.events = simulator._generate_events()
    
    if args.customers:
        simulator.config['simulation']['customers']['count'] = args.customers
        simulator.customers = simulator._generate_customers()
    
    if args.duration:
        simulator.config['simulation']['duration_minutes'] = args.duration
    
    if args.purchases_per_min:
        simulator.config['simulation']['traffic']['purchases_per_minute'] = args.purchases_per_min
    
    # Run simulation
    simulator.run_simulation(scenario=args.scenario)


if __name__ == '__main__':
    main()

Create simulators/README.md:

# Ticket Marketplace Simulator

Configurable event simulator for ticket marketplace analytics.

## Quick Start

```bash
# Install dependencies
pip install -r requirements.txt

# Run default simulation
python ticket_marketplace_simulator.py

# Run with specific scenario
python ticket_marketplace_simulator.py --scenario flash_sale

# Run with custom configuration
python ticket_marketplace_simulator.py --vendors 100 --customers 5000 --duration 30

# Run with custom purchase rate
python ticket_marketplace_simulator.py --purchases-per-min 500 --duration 10

Available Scenarios

  • normal_day - Normal traffic (50 purchases/min, 60 minutes)
  • flash_sale - High traffic flash sale (500 purchases/min, 15 minutes)
  • concert_announcement - Concert announcement spike (200 purchases/min, 30 minutes)
  • weekend_rush - Weekend traffic (300 purchases/min, 120 minutes)

Configuration

Edit config.yaml to customize:

  • Number of vendors, customers, events
  • Traffic rates (purchases, views, cart events, searches)
  • Peak hours and multipliers
  • Conversion funnel rates
  • Event categories and ticket types
  • Payment method distribution

Command Line Options

--config CONFIG         Configuration file path (default: config.yaml)
--scenario SCENARIO     Scenario name to run
--vendors N            Override number of vendors
--customers N          Override number of customers
--duration N           Override duration in minutes
--purchases-per-min N  Override purchases per minute

Examples

# Black Friday simulation - 10 minutes of intense traffic
python ticket_marketplace_simulator.py --purchases-per-min 1000 --duration 10

# Long-running stress test - 2 hours
python ticket_marketplace_simulator.py --duration 120 --vendors 200 --customers 10000

# Low traffic test
python ticket_marketplace_simulator.py --purchases-per-min 10 --duration 5

---

## Step 6: Environment Variables

Create `.env`:

```env
# Kafka
KAFKA_BOOTSTRAP_SERVERS=localhost:9092

# ClickHouse
CLICKHOUSE_HOST=localhost
CLICKHOUSE_PORT=8123
CLICKHOUSE_USER=default
CLICKHOUSE_PASSWORD=
CLICKHOUSE_DATABASE=ticket_analytics

# dbt
DBT_PROFILES_DIR=./dbt

# Metabase
MB_DB_FILE=/metabase-data/metabase.db

Usage Instructions

1. Start All Services

# Navigate to project directory
cd ticket-analytics-engine

# Start all containers
docker-compose up -d

# Check status
docker-compose ps

# View logs
docker-compose logs -f

2. Verify Services

3. Create Kafka Topics

# Create all required topics
docker exec -it kafka kafka-topics \
  --create \
  --topic ticket-purchases \
  --bootstrap-server localhost:9092 \
  --partitions 3 \
  --replication-factor 1

docker exec -it kafka kafka-topics \
  --create \
  --topic ticket-views \
  --bootstrap-server localhost:9092 \
  --partitions 3 \
  --replication-factor 1

docker exec -it kafka kafka-topics \
  --create \
  --topic cart-events \
  --bootstrap-server localhost:9092 \
  --partitions 3 \
  --replication-factor 1

docker exec -it kafka kafka-topics \
  --create \
  --topic search-events \
  --bootstrap-server localhost:9092 \
  --partitions 3 \
  --replication-factor 1

# List all topics
docker exec -it kafka kafka-topics --list --bootstrap-server localhost:9092

4. Seed Initial Data (Optional)

Load vendors, customers, and events into ClickHouse:

# Execute SQL to insert seed data
docker exec -it clickhouse clickhouse-client --query "
INSERT INTO ticket_analytics.vendors VALUES
('vendor_0001', 'Live Nation', 'Music Concerts', 'USA', 'Los Angeles', 10.5, 'active', now(), now()),
('vendor_0002', 'Ticketmaster Sports', 'Sports', 'USA', 'New York', 12.0, 'active', now(), now());
-- Add more vendors...
"

5. Run Ticket Marketplace Simulator

cd simulators
pip install -r requirements.txt

# Run default simulation (10 minutes)
python ticket_marketplace_simulator.py

# Run flash sale scenario (high traffic)
python ticket_marketplace_simulator.py --scenario flash_sale

# Custom configuration
python ticket_marketplace_simulator.py \
  --vendors 100 \
  --customers 5000 \
  --duration 30 \
  --purchases-per-min 200

# Weekend rush (2 hours)
python ticket_marketplace_simulator.py --scenario weekend_rush

6. Monitor Real-Time Data

# Monitor Kafka topics
docker exec -it kafka kafka-console-consumer \
  --bootstrap-server localhost:9092 \
  --topic ticket-purchases \
  --from-beginning \
  --max-messages 10

# Check ClickHouse data ingestion
docker exec -it clickhouse clickhouse-client --query "
SELECT count() as total_purchases 
FROM ticket_analytics.ticket_purchases_raw;
"

# View real-time aggregations
docker exec -it clickhouse clickhouse-client --query "
SELECT 
    toStartOfHour(purchase_timestamp) as hour,
    count() as purchases,
    sum(total_price) as revenue
FROM ticket_analytics.ticket_purchases_raw
GROUP BY hour
ORDER BY hour DESC
LIMIT 24;
"

7. Run dbt Transformations

cd dbt

# Install dbt-clickhouse
pip install dbt-clickhouse

# Test connection
dbt debug

# Run all models
dbt run

# Run specific model
dbt run --select vendor_performance

# Run tests
dbt test

# Generate documentation
dbt docs generate
dbt docs serve  # Opens at http://localhost:8080

8. Configure Metabase Dashboards

  1. Navigate to http://localhost:3000

  2. Complete initial setup

  3. Add ClickHouse database:

    • Database type: ClickHouse (requires plugin or use generic database with JDBC)
    • Host: clickhouse
    • Port: 8123
    • Database: ticket_analytics
    • Username: default
    • Password: (leave empty)
  4. Create dashboards with these queries:

Vendor Performance Dashboard:

SELECT * FROM ticket_analytics.vendor_performance
ORDER BY total_revenue DESC
LIMIT 20

Hourly Sales Trend:

SELECT 
    metric_hour,
    total_transactions,
    total_revenue,
    unique_customers
FROM ticket_analytics.marketplace_hourly_metrics
ORDER BY metric_hour DESC
LIMIT 48

Top Events:

SELECT * FROM ticket_analytics.ticket_sales_metrics
ORDER BY revenue DESC
LIMIT 10

Customer Segmentation:

SELECT 
    customer_segment,
    COUNT(*) as customer_count,
    SUM(total_spent) as segment_revenue,
    AVG(total_spent) as avg_customer_value
FROM ticket_analytics.customer_behavior
GROUP BY customer_segment

Simulation Scenarios

Normal Day

python ticket_marketplace_simulator.py --scenario normal_day
  • 50 purchases/minute
  • 300 views/minute
  • 60 minutes duration
  • Simulates typical weekday traffic

Flash Sale

python ticket_marketplace_simulator.py --scenario flash_sale
  • 500 purchases/minute
  • 2000 views/minute
  • 15 minutes duration
  • 5x peak multiplier
  • Simulates limited-time offers

Concert Announcement

python ticket_marketplace_simulator.py --scenario concert_announcement
  • 200 purchases/minute
  • 1500 views/minute
  • 30 minutes duration
  • Simulates major event announcement spike

Weekend Rush

python ticket_marketplace_simulator.py --scenario weekend_rush
  • 300 purchases/minute
  • 1200 views/minute
  • 120 minutes duration
  • 2.5x peak multiplier

Custom Scenarios

Create your own scenario in config.yaml:

scenarios:
  my_custom_scenario:
    purchases_per_minute: 150
    views_per_minute: 800
    cart_events_per_minute: 250
    searches_per_minute: 200
    duration_minutes: 45
    peak_multiplier: 2.0

Run it:

python ticket_marketplace_simulator.py --scenario my_custom_scenario

Key Metrics to Monitor

Vendor Metrics

  • Total revenue per vendor
  • Tickets sold
  • Unique customers
  • Conversion rate (views to purchases)
  • Average transaction value
  • Events with sales

Event Metrics

  • Sell-through rate
  • Revenue per ticket
  • Unique buyers
  • Sales duration
  • Tickets remaining

Customer Metrics

  • Customer lifetime value
  • Purchase frequency
  • Average transaction value
  • Customer segments (VIP, High Value, Medium, Low)
  • Activity status (Active, At Risk, Churned)
  • View-to-purchase conversion rate

Marketplace Metrics

  • Hourly transaction volume
  • Total GMV (Gross Merchandise Value)
  • Unique active customers
  • Active vendors
  • Average order value
  • Conversion funnel metrics

GDPR Compliance Features

Data Retention

Automatically delete old data:

-- Delete events older than 90 days (run as scheduled job)
ALTER TABLE ticket_analytics.ticket_purchases_raw 
DELETE WHERE purchase_timestamp < now() - INTERVAL 90 DAY;

ALTER TABLE ticket_analytics.ticket_views_raw 
DELETE WHERE view_timestamp < now() - INTERVAL 90 DAY;

ALTER TABLE ticket_analytics.cart_events_raw 
DELETE WHERE event_timestamp < now() - INTERVAL 90 DAY;

Right to be Forgotten

-- Mark customer for deletion
INSERT INTO ticket_analytics.user_deletions (customer_id, deletion_requested_at)
VALUES ('customer_00123', now());

-- Delete customer data from all tables
ALTER TABLE ticket_analytics.ticket_purchases_raw 
DELETE WHERE customer_id = 'customer_00123';

ALTER TABLE ticket_analytics.ticket_views_raw 
DELETE WHERE customer_id = 'customer_00123';

ALTER TABLE ticket_analytics.cart_events_raw 
DELETE WHERE customer_id = 'customer_00123';

ALTER TABLE ticket_analytics.search_events_raw 
DELETE WHERE customer_id = 'customer_00123';

-- Update deletion status
ALTER TABLE ticket_analytics.user_deletions 
UPDATE deletion_completed_at = now(), deletion_status = 'completed'
WHERE customer_id = 'customer_00123';

Data Export (Right to Access)

-- Export all customer data
SELECT 
    'purchases' as data_type,
    * 
FROM ticket_analytics.ticket_purchases_raw
WHERE customer_id = 'customer_00123'
UNION ALL
SELECT 
    'views' as data_type,
    view_id, customer_id, vendor_id, event_id, 
    view_timestamp, view_duration_seconds, session_id, user_consent, created_at
FROM ticket_analytics.ticket_views_raw
WHERE customer_id = 'customer_00123'
FORMAT JSONEachRow;

Consent Management

-- Update customer consent status
ALTER TABLE ticket_analytics.customers 
UPDATE user_consent = false 
WHERE customer_id = 'customer_00123';

-- Only consented data is processed (enforced in Flink jobs)
SELECT * FROM ticket_analytics.ticket_purchases_raw
WHERE user_consent = true;

Scaling Horizontally

Scale Kafka Brokers

Add more brokers in docker-compose.yml:

kafka-2:
  image: confluentinc/cp-kafka:7.5.0
  hostname: kafka-2
  container_name: kafka-2
  depends_on:
    - zookeeper
  ports:
    - "9094:9094"
  environment:
    KAFKA_BROKER_ID: 2
    KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
    KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
    KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka-2:29093,PLAINTEXT_HOST://localhost:9094
    # ... other settings

Scale Flink TaskManagers

# Scale to 4 task managers
docker-compose up -d --scale flink-taskmanager=4

# Verify
docker-compose ps flink-taskmanager

Scale ClickHouse (Cluster Setup)

For production, set up ClickHouse cluster with replication:

  1. Add more ClickHouse nodes
  2. Configure distributed tables
  3. Set up ZooKeeper for coordination

Example cluster configuration:

<!-- config.xml -->
<remote_servers>
    <ticket_analytics_cluster>
        <shard>
            <replica>
                <host>clickhouse-1</host>
                <port>9000</port>
            </replica>
            <replica>
                <host>clickhouse-2</host>
                <port>9000</port>
            </replica>
        </shard>
        <shard>
            <replica>
                <host>clickhouse-3</host>
                <port>9000</port>
            </replica>
            <replica>
                <host>clickhouse-4</host>
                <port>9000</port>
            </replica>
        </shard>
    </ticket_analytics_cluster>
</remote_servers>

Performance Optimization

ClickHouse Optimizations

-- Optimize table after bulk inserts
OPTIMIZE TABLE ticket_analytics.ticket_purchases_raw FINAL;

-- Create secondary indices for common queries
ALTER TABLE ticket_analytics.ticket_purchases_raw 
ADD INDEX idx_customer_id customer_id TYPE bloom_filter GRANULARITY 1;

ALTER TABLE ticket_analytics.ticket_purchases_raw 
ADD INDEX idx_vendor_id vendor_id TYPE bloom_filter GRANULARITY 1;

-- Materialize computed columns
ALTER TABLE ticket_analytics.ticket_purchases_raw 
ADD COLUMN purchase_date Date MATERIALIZED toDate(purchase_timestamp);

Kafka Optimizations

# Increase partitions for higher throughput
docker exec -it kafka kafka-topics \
  --alter \
  --topic ticket-purchases \
  --partitions 6 \
  --bootstrap-server localhost:9092

Flink Optimizations

  • Increase parallelism for Flink jobs
  • Tune checkpoint intervals
  • Optimize buffer sizes
  • Use RocksDB state backend for larger states

Monitoring & Troubleshooting

Check ClickHouse Performance

-- Query statistics
SELECT 
    type,
    query_kind,
    event_time,
    query_duration_ms,
    read_rows,
    read_bytes,
    memory_usage
FROM system.query_log 
WHERE type = 'QueryFinish' 
  AND event_time >= now() - INTERVAL 1 HOUR
ORDER BY query_duration_ms DESC 
LIMIT 10;

-- Table sizes
SELECT 
    table,
    formatReadableSize(sum(bytes)) as size,
    sum(rows) as rows,
    max(modification_time) as latest_modification
FROM system.parts
WHERE database = 'ticket_analytics' AND active
GROUP BY table
ORDER BY sum(bytes) DESC;

-- Partition info
SELECT 
    partition,
    sum(rows) as rows,
    formatReadableSize(sum(bytes)) as size
FROM system.parts
WHERE database = 'ticket_analytics' 
  AND table = 'ticket_purchases_raw'
  AND active
GROUP BY partition
ORDER BY partition DESC;

Kafka Monitoring

# Consumer lag
docker exec -it kafka kafka-consumer-groups \
  --bootstrap-server localhost:9092 \
  --describe \
  --group flink-purchases-consumer

# Topic details
docker exec -it kafka kafka-topics \
  --describe \
  --topic ticket-purchases \
  --bootstrap-server localhost:9092

# Message count (approximate)
docker exec -it kafka kafka-run-class kafka.tools.GetOffsetShell \
  --broker-list localhost:9092 \
  --topic ticket-purchases

Flink Monitoring

  • Access Flink UI: http://localhost:8081
  • Check job status, metrics, and checkpoints
  • View task manager resources
  • Monitor backpressure
# View Flink logs
docker logs flink-jobmanager
docker logs flink-taskmanager

Simulator Statistics

The simulator provides real-time statistics:

⏳ Elapsed: 5.2min | Purchases: 1042 | Views: 2613 | Cart: 1020 | Searches: 781

Cleanup

# Stop all services
docker-compose down

# Remove volumes (WARNING: deletes all data)
docker-compose down -v

# Remove images
docker-compose down --rmi all

# Clean everything including orphaned containers
docker-compose down -v --remove-orphans

Next Steps

  1. Add Real-time Alerting:

    • Integrate with Prometheus + Grafana
    • Set up alerts for anomalies (sudden drops in purchases, high error rates)
  2. Implement Data Quality Checks:

    • Add Great Expectations or dbt tests
    • Monitor data freshness and completeness
  3. Add Machine Learning:

    • Build recommendation engine
    • Predict event popularity
    • Customer churn prediction
    • Dynamic pricing models
  4. Setup CI/CD:

    • Automate dbt runs
    • Deploy Flink jobs automatically
    • Infrastructure as Code (Terraform)
  5. Add More Data Sources:

    • Payment gateway webhooks
    • Email marketing events
    • Customer support interactions
    • Social media sentiment
  6. Optimize for Production:

    • Setup authentication and authorization
    • Implement rate limiting
    • Add data encryption
    • Setup backup and disaster recovery
  7. Advanced Analytics:

    • Cohort analysis
    • Attribution modeling
    • Customer lifetime value prediction
    • Market basket analysis

Resources


License

MIT License - Feel free to modify and use for your projects.