Metadata
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
Application entry point. Constructs the CLI application, registers all subcommands and flags, and invokes app.Run(os.Args).
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