Jump to content

Connect SuperML | Leeroopedia MCP: Equip your AI agents with best practices, code verification, and debugging knowledge. Powered by Leeroo — building Organizational Superintelligence. Contact us at founders@leeroo.com.

Implementation:Risingwavelabs Risingwave Datagen CLI

From Leeroopedia


Metadata

Property Value
File integration_tests/datagen/main.go
Package main
Language Go
Lines 283
Category Integration Test CLI Entry Point
Repository https://github.com/risingwavelabs/risingwave

Overview

The Datagen CLI is the main executable entry point for the datagen tool, which is used by RisingWave's integration test docker-compose setups to generate synthetic data. It uses the urfave/cli library to define a command-line interface with subcommands for each supported sink type and global flags for controlling the data generation behavior.

The CLI defines seven sink subcommands (postgres, mysql, kafka, pulsar, kinesis, s3, nats), each with sink-specific flags. Seven global flags control cross-cutting concerns such as the generation mode, output format, QPS throttling, topic filtering, and total event count.

When invoked, the CLI parses command-line arguments into a gen.GeneratorConfig struct, sets up signal handling for graceful shutdown (SIGINT, SIGTERM), and delegates to the generateLoad function to execute the data generation loop.

Code Reference

Source Location

integration_tests/datagen/main.go

Signature

func main()

Application entry point. Constructs the CLI application, registers all subcommands and flags, and invokes app.Run(os.Args).

func runCommand() error

Sets up signal handling and a cancellable context, then calls generateLoad(ctx, cfg).

Import

import (
	"context"
	"datagen/gen"
	"log"
	"os"
	"os/signal"
	"syscall"

	"github.com/urfave/cli"
)

I/O Contract

Global Flags

Flag Type Required Default Description
--mode string Yes -- Data generation mode: ad-click, ad-ctr, twitter, cdn-metrics, clickstream, ecommerce, delivery, livestream, compatible-data
--qps int No 1 Number of messages to send per second
--format string No "json" Output record format: json or protobuf (used for message queue sinks)
--print bool No false Whether to print every event's SQL representation to stdout
--heavytail bool No false Use uniform distribution for randomizing values (higher tail probability)
--topic string No "" Topic filter; if set, only records matching this topic are emitted
--total_event int64 No 0 Total events to generate; 0 means run indefinitely

Sink Subcommands

postgres

Flag Type Default Description
--host string "localhost" PostgreSQL server host address
--db string "dev" Target database name
--port int 4566 PostgreSQL server port
--user string "root" PostgreSQL user

mysql

Flag Type Default Description
--host string "localhost" MySQL server host address
--db string "mydb" Target database name
--port int 3306 MySQL server port
--user string "mysqluser" MySQL user
--password string "mysqlpw" MySQL password

kafka

Flag Type Required Description
--brokers string Yes Comma-separated list of Kafka bootstrap brokers
--no-recreate bool No Do not recreate Kafka topics if they already exist

pulsar

Flag Type Required Description
--brokers string Yes Comma-separated list of Pulsar brokers

kinesis

Flag Type Required Description
--region string Yes AWS region of the Kinesis stream
--endpoint string No Kinesis stream endpoint
--name string Yes Kinesis stream name

s3

Flag Type Required Description
--region string Yes AWS region of the S3 bucket
--bucket string Yes S3 bucket name
--endpoint string No S3 bucket endpoint

nats

Flag Type Required Description
--url string Yes URL of the NATS server
--jetstream bool No Whether to use JetStream

Signal Handling

func runCommand() error {
	terminateCh := make(chan os.Signal, 1)
	signal.Notify(terminateCh, os.Interrupt, syscall.SIGTERM)

	ctx, cancel := context.WithCancel(context.Background())
	go func() {
		<-terminateCh
		log.Println("Cancelled")
		cancel()
	}()
	return generateLoad(ctx, cfg)
}

The runCommand function listens for SIGINT and SIGTERM, cancelling the context to trigger a graceful shutdown of the generation loop.

Usage Examples

Generating Ad CTR Data to Kafka

// Command line:
// datagen --mode ad-ctr --qps 100 --format json kafka --brokers localhost:9092

Generating Ecommerce Data to PostgreSQL

// Command line:
// datagen --mode ecommerce --qps 50 --print postgres --host localhost --port 4566 --db dev --user root

Generating Limited Events to S3

// Command line:
// datagen --mode clickstream --total_event 10000 --format json s3 --region us-east-1 --bucket my-bucket

Filtering by Topic

// Only generate ad_impression events (skip ad_click):
// datagen --mode ad-ctr --topic ad_impression kafka --brokers localhost:9092

Related Pages

Page Connections

Double-click a node to navigate. Hold to expand connections.
Principle
Implementation
Heuristic
Environment