Skip to content

Latest commit

 

History

1,685 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Nisshi (née Tansu) 🗃️

stateless Kafka-compatible broker with pluggable storage (PostgreSQL, SQLite, S3, memory)


License   Built with Rust  
Docs   Blog  
GitHub stars   Bluesky


What is Nisshi?

Nisshi (née Tansu) is a drop-in replacement for Apache Kafka with PostgreSQL, libSQL (SQLite), S3 or memory storage engines. Schema backed topics (Avro, JSON or Protocol buffers) can be written as Apache Iceberg or Delta Lake tables.

Features:

See examples using pyiceberg, examples using Apache Spark or 🆕 examples using Delta Lake.

For data durability:

Nisshi is a single statically linked binary containing the following:

  • broker an Apache Kafka API compatible broker and schema registry
  • topic a CLI to create/delete Topics
  • cat a CLI to consume or produce Avro, JSON or Protobuf messages to a topic
  • proxy an Apache Kafka compatible proxy

broker

The broker subcommand is default if no other command is supplied.

Usage: nisshi [OPTIONS]
       nisshi <COMMAND>

Commands:
  broker  Apache Kafka compatible broker with Avro, JSON, Protobuf schema validation [default if no command supplied]
  cat     Easily consume or produce Avro, JSON or Protobuf messages to a topic
  topic   Create or delete topics managed by the broker
  proxy   Apache Kafka compatible proxy
  help    Print this message or the help of the given subcommand(s)

Options:
      --kafka-cluster-id <KAFKA_CLUSTER_ID>
          All members of the same cluster should use the same id [env: CLUSTER_ID=RvQwrYegSUCkIPkaiAZQlQ] [default: nisshi_cluster]
      --kafka-listener-url <KAFKA_LISTENER_URL>
          The broker will listen on this address [env: LISTENER_URL=] [default: tcp://[::]:9092]
      --kafka-advertised-listener-url <KAFKA_ADVERTISED_LISTENER_URL>
          This location is advertised to clients in metadata [env: ADVERTISED_LISTENER_URL=tcp://localhost:9092] [default: tcp://localhost:9092]
      --storage-engine <STORAGE_ENGINE>
          Storage engine examples are: postgres://postgres:postgres@localhost, memory://nisshi/ or s3://nisshi/ [env: STORAGE_ENGINE=s3://nisshi/] [default: memory://nisshi/]
      --schema-registry <SCHEMA_REGISTRY>
          Schema registry examples are: file://./etc/schema or s3://nisshi/, containing: topic.json, topic.proto or topic.avsc [env: SCHEMA_REGISTRY=file://./etc/schema]
      --data-lake <DATA_LAKE>
          Apache Parquet files are written to this location, examples are: file://./lake or s3://lake/ [env: DATA_LAKE=s3://lake/]
      --iceberg-catalog <ICEBERG_CATALOG>
          Apache Iceberg Catalog, examples are: http://localhost:8181/ [env: ICEBERG_CATALOG=http://localhost:8181/]
      --iceberg-namespace <ICEBERG_NAMESPACE>
          Iceberg namespace [env: ICEBERG_NAMESPACE=] [default: nisshi]
      --prometheus-listener-url <PROMETHEUS_LISTENER_URL>
          Broker metrics can be scraped by Prometheus from this URL [env: PROMETHEUS_LISTENER_URL=tcp://0.0.0.0:9100] [default: tcp://[::]:9100]
  -h, --help
          Print help
  -V, --version
          Print version

A broker can be started by simply running nisshi, all options have defaults. Nisshi pickup any existing environment, loading any found in .env. An example.env is provided as part of the distribution and can be copied into .env for local modification. Sample schemas can be found in etc/schema, used in the examples.

If an Apache Avro, Protobuf or JSON schema has been assigned to a topic, the broker will reject any messages that are invalid. Schema backed topics are written as Apache Parquet when the -data-lake option is provided.

TLS

The listener speaks TLS when both --cert (the certificate chain) and --key (the private key) are supplied as PEM files, and clap rejects one without the other:

nisshi broker --cert broker.pem --key broker-key.pem

With TLS configured the listener is TLS only: plaintext clients are refused during the handshake, so clients must be configured with security.protocol=SSL (and a truststore containing the certificate, if it is self-signed). The --listener-url keeps its tcp:// scheme. TLS here provides encryption only; client authentication is still SASL (see --authentication). The private key is a PKCS#8, SEC1 or RSA PEM key. --cert and --key may point at the same file when the certificate chain and key are kept in one PEM bundle. Any problem loading the certificate or key, or a key that does not match the certificate, fails startup rather than falling back to plaintext.

A passphrase protected key (the counterpart of Kafka's ssl.key.password) is read with --key-passphrase-file, a file holding the passphrase; a trailing newline is ignored:

nisshi broker --cert broker.pem --key broker-key.pem --key-passphrase-file broker-key.passphrase

Only PKCS#8 encryption (ENCRYPTED PRIVATE KEY, PBES2 with PBKDF2-HMAC-SHA2 or scrypt and AES-CBC or Triple DES, which is what openssl pkcs8 -topk8 emits) is supported. Legacy OpenSSL PEM encryption (Proc-Type: 4,ENCRYPTED) and keys derived with a SHA-1 PRF (older OpenSSL releases' default) are rejected at startup, naming the algorithm; re-encrypt the key first, keeping the passphrase:

openssl pkcs8 -topk8 -in broker-key.pem -out broker-key-pkcs8.pem -v2 aes-256-cbc -v2prf hmacWithSHA256

An empty passphrase file means no passphrase, so a mounted secret that is absent for unencrypted keys works unchanged. A passphrase given for a key that is not encrypted is ignored with a warning.

Note for existing deployments: before 0.7 these flags were accepted but had no effect, so a broker started with --cert and --key was serving plaintext. After upgrading, that same command line serves TLS only, and --cert without --key (or the reverse) is rejected. The bundled nisshi cat, topic, perf and proxy subcommands connect in plaintext and cannot yet talk to a TLS listener. --key-passphrase-file is unknown to earlier releases, so rolling the binary back means removing that flag from the command line as well.

Metrics

Metrics are exported over OTLP/HTTP when --otlp-endpoint-url (or OTEL_EXPORTER_OTLP_ENDPOINT) is set. v1/metrics is appended to the URL's path, so https://collector:4318 and https://gateway/otlp export to https://collector:4318/v1/metrics and https://gateway/otlp/v1/metrics. OTEL_EXPORTER_OTLP_METRICS_ENDPOINT is not read. Collector credentials go in OTEL_EXPORTER_OTLP_HEADERS.

The exported resource honours the standard OpenTelemetry environment variables. Every key=value pair in OTEL_RESOURCE_ATTRIBUTES is attached as given; values are not percent-decoded, and entries without a = are ignored with a warning. service.name is resolved as OTEL_SERVICE_NAME, then service.name in OTEL_RESOURCE_ATTRIBUTES, then the default for the subcommand (nisshi-broker, or nisshi-proxy / nisshi-generator for proxy and generator). Unset, empty or whitespace-only values fall through to the next source. The endpoint, service.name and attribute names are logged at info on startup. Resource attributes are sent with every export, so use an https endpoint when they carry anything sensitive.

OTEL_EXPORTER_OTLP_ENDPOINT=https://collector:4318/ \
OTEL_SERVICE_NAME=kafka-broker \
OTEL_RESOURCE_ATTRIBUTES=service.version=0.7.0,deployment.environment.name=staging \
nisshi broker

Note for existing deployments: before 0.7 service.name was always nisshi-broker and these variables were ignored. If OTEL_SERVICE_NAME or OTEL_RESOURCE_ATTRIBUTES is already set in the broker's environment, for example injected pod-wide for another agent, its series move to that name on upgrade and dashboards or alerts keyed on nisshi-broker go quiet. Unset the inherited variable or set OTEL_SERVICE_NAME=nisshi-broker. Rolling back reverts to nisshi-broker regardless of these variables. On shutdown the broker now flushes pending metrics before exiting.

topic

The nisshi topic command has the following subcommands:

Create or delete topics managed by the broker

Usage: nisshi topic <COMMAND>

Commands:
  create  Create a topic
  delete  Delete an existing topic
  help    Print this message or the help of the given subcommand(s)

Options:
  -h, --help  Print help

To create a topic use:

nisshi topic create taxi

cat

The nisshi cat command, has the following subcommands:

nisshi cat --help
Easily consume or produce Avro, JSON or Protobuf messages to a topic

Usage: nisshi cat <COMMAND>

Commands:
  produce  Produce Avro/JSON/Protobuf messages to a topic
  consume  Consume Avro/JSON/Protobuf messages from a topic
  help     Print this message or the help of the given subcommand(s)

Options:
  -h, --help  Print help

The produce subcommand reads JSON formatted messages encoding them into Apache Avro, Protobuf or JSON depending on the schema used by the topic.

For example, the taxi topic is backed by taxi.proto. Using trips.json containing a JSON array of objects, nisshi cat produce encodes each message into protobuf into the broker:

nisshi cat produce taxi etc/data/trips.json

Using duckdb we can read the Apache Parquet files created by the broker:

duckdb :memory: "SELECT * FROM 'data/taxi/*/*.parquet'"

Results in the following output:

|-----------+---------+---------------+-------------+---------------|
| vendor_id | trip_id | trip_distance | fare_amount | store_and_fwd |
|     int64 |   int64 |         float |      double |         int32 |
|-----------+---------+---------------+-------------+---------------|
|         1 | 1000371 |           1.8 |       15.32 |             0 |
|         2 | 1000372 |           2.5 |       22.15 |             0 |
|         2 | 1000373 |           0.9 |        9.01 |             0 |
|         1 | 1000374 |           8.4 |       42.13 |             1 |
|-----------+---------+---------------+-------------+---------------|

s3

The following will configure a S3 storage engine using the "nisshi" bucket (full context is in compose.yaml and example.env):

Copy example.env into .env so that you have a local working copy:

cp example.env .env

Edit .env so that STORAGE_ENGINE is defined as:

STORAGE_ENGINE="s3://nisshi/"

First time startup, you'll need to create a bucket, an access key and a secret in minio.

Just bring minio up, without nisshi:

docker compose up -d minio

Create a minio local alias representing http://localhost:9000 with the default credentials of minioadmin:

docker compose exec minio \
   /usr/bin/mc \
   alias \
   set \
   local \
   http://localhost:9000 \
   minioadmin \
   minioadmin

Create a nisshi bucket in minio using the local alias:

docker compose exec minio \
   /usr/bin/mc mb local/nisshi

Once this is done, you can start nisshi with:

docker compose up -d nisshi

Using the regular Apache Kafka CLI you can create topics, produce and consume messages with Nisshi:

kafka-topics \
  --bootstrap-server localhost:9092 \
  --partitions=3 \
  --replication-factor=1 \
  --create --topic test

Describe the test topic:

kafka-topics \
  --bootstrap-server localhost:9092 \
  --describe \
  --topic test

Note that node 111 is the leader and ISR for each topic partition. This node represents the broker handling your request. All brokers are node 111.

Producer:

echo "hello world" | kafka-console-producer \
    --bootstrap-server localhost:9092 \
    --topic test

Group consumer using test-consumer-group:

kafka-console-consumer \
  --bootstrap-server localhost:9092 \
  --group test-consumer-group \
  --topic test \
  --from-beginning \
  --property print.timestamp=true \
  --property print.key=true \
  --property print.offset=true \
  --property print.partition=true \
  --property print.headers=true \
  --property print.value=true

Describe the consumer test-consumer-group group:

kafka-consumer-groups \
  --bootstrap-server localhost:9092 \
  --group test-consumer-group \
  --describe

PostgreSQL

To switch between the minio and PostgreSQL examples, firstly shutdown Nisshi:

docker compose down nisshi

Switch to the PostgreSQL storage engine by updating .env:

# minio storage engine
# STORAGE_ENGINE="s3://nisshi/"

# PostgreSQL storage engine -- NB: @db and NOT @localhost :)
STORAGE_ENGINE="postgres://postgres:postgres@db"

Start PostgreSQL:

docker compose up -d db

Bring Nisshi back up:

docker compose up -d nisshi

Using the regular Apache Kafka CLI you can create topics, produce and consume messages with Nisshi:

kafka-topics \
  --bootstrap-server localhost:9092 \
  --partitions=3 \
  --replication-factor=1 \
  --create --topic test

Producer:

echo "hello world" | kafka-console-producer \
    --bootstrap-server localhost:9092 \
    --topic test

Consumer:

kafka-console-consumer \
  --bootstrap-server localhost:9092 \
  --group test-consumer-group \
  --topic test \
  --from-beginning \
  --property print.timestamp=true \
  --property print.key=true \
  --property print.offset=true \
  --property print.partition=true \
  --property print.headers=true \
  --property print.value=true

Or using librdkafka to produce:

echo "Lorem ipsum dolor..." | \
  ./examples/rdkafka_example -P \
  -t test -p 1 \
  -b localhost:9092 \
  -z gzip

Consumer:

./examples/rdkafka_example \
  -C \
  -t test -p 1 \
  -b localhost:9092

Feedback

Please raise an issue if you encounter a problem.

Contributing

Contributions are welcome. See CONTRIBUTING for how to get started, and GOVERNANCE for how decisions get made. Everyone participating is expected to follow our Code of Conduct. To report a security vulnerability, see SECURITY.

License

Nisshi is licensed under Apache 2.0.

About

Apache Kafka® compatible broker with S3, PostgreSQL, SQLite, Apache Iceberg and Delta Lake

Topics

Resources

Code of conduct

Contributing

Security policy

Stars

1.9k stars

Watchers

10 watching

Forks

Releases

Packages

Used by

Contributors

Languages