1. User Manual
1.1. Introduction
This user manual will help you understand kafka Cost Control and how to use it propertly. This document assumes that you already have a running application. If not please see the Installation section.
At this point you should have access to the Kafka Cost Control UI.
1.1.1. Graphql
Kafka cost control provides a graphql endpoint at: <your-host>/graphql-ui
In addition, there is a ready to use GraphQL UI. You can access it by going to the following URL: <your-host>/graphql-ui
1.1.2. Authorization
TODO explain basic auth stuff TODO explain the localstorage trick
1.2. Config samples
If you want to quickly get started, you can create pricing rules and context data using the config sample folder. All you need is node JS 20+.
Be sure to edit the files pricing-rules.json and context-data.json to match your environment.
To persist the configuration you can use the following command:
node index.js --url=https://<your-host>/graphql --user=admin --password=your-password
If you have issues you can try to add the --verbose options. This will display all the requests.
1.3. Configuring your cost control application
Your cost control instance can be configured either by setting environment variables,
or (if you are using one of the provided container images) by mounting an application.properties file at
/deployments/config/application.properties:
# inside your application.properties
quarkus.profile=ccloud
1.3.1. Configuring aggregation types
In cost control, each metric is associated with an aggregation type. The aggregation type determines how multiple measurements
within a time window are combined to produce a single value. Possible aggregation types are SUM and MAX.
If no aggregation type is specified for a metric, then SUM will be used.
To specify an aggregation type for a metric, add a line to your application.properties file like this:
# inside your application.properties
cc.metrics.aggregations.<metric-name>=max
For example, to aggregate retained byte measurements using the MAX aggregation type, add the following line to your application.properties file:
cc.metrics.aggregations.confluent_kafka_server_retained_bytes=max
Note that these settings may only be set in an application.properties file, and not via environment variables.
This is due to the fact that Quarkus cannot unambiguously map environment variable names to property names if said
property names contain user-defined parts (like confluent_kafka_server_retained_bytes in the example above).
For more information, see the relevant section of the Quarkus documentation.
1.3.2. Configuring transformations
Currently, there is one kind of transformation that can be added to the metrics processing pipeline: splitTopicMetricAmongPrincipals.
This transformation is configured by supplying a map of metric names to context keys in application.properties. For example:
cc.metrics.transformations.splitMetricAmongPrincipals={bytesin:'writers',bytesout:'readers'}
This will expect bytesin metrics to have a writers key in their context, which should be a comma-separated list of principals (e.g. applications, teams, …)
that can write to this topic. When a bytesin metric for any topic is encountered, the metric will be replaced with n metrics (where n is the length of the writers list
in that metric’s context). The name of the topic will be moved into that metric’s context under the topic key and the name of a principal from the writers
list will move into the metric’s name field. The value of each generated metric will be the value of the original metric divided by n.
In effect this takes a single metric that says "90 kb have been written to topic ABC by writers X,Y,Z" and transforms it into 3 metrics:
-
Principal X wrote 30 kb to topic ABC
-
Principal Y wrote 30 kb to topic ABC
-
Principal Z wrote 30 kb to topic ABC
The advantage over the original metric is that each such transformed metric will be put into its own database row, which will make it easier to aggregate by principal (e.g. answer questions like "how much did principal X produce in the past month"). This type of transformation is especially relevant when doing cost control based on vanilla Kafka broker metrics (as is the case in Strimzi deployments, for example), because here only topic-level metrics are generated natively. This transformation allows to turn these topic-focused metrics into principal-focused ones.
Some metrics that are configured this way may not have a matching context key in the context. To handle such cases, cost control offers three different strategies
that can be enabled by setting the cc.metrics.transformations.config.splitMetricAmongPrincipals.missingKeyHandling property in application.properties:
-
ASSIGN_TO_FALLBACK: The metric will be handled as if the context key existed and contained a single principal name. This fallback principal name is set to "unknown" per default but can be changed by setting themetrics.transformations.config.splitMetricAmongPrincipals.fallbackPrincipalproperty inapplication.properties. -
DROP: The metric will not be forwarded downstream and will be dropped. -
PASS_THROUGH: The metric will be passed downstream without any changes. This is the default behavior.
Note that cc.metrics.transformations.config.splitMetricAmongPrincipals.missingKeyHandling is expected to be a map of metric names to one of the above strategies.
For example:
cc.metrics.transformations.config.splitMetricAmongPrincipals.missingKeyHandling={bytesin:'ASSIGN_TO_FALLBACK',bytesout:'DROP'}
If no strategy is specified for a metric, then ASSIGN_TO_FALLBACK will be used.
1.4. Pricing rules
Pricing rules are a way to put a price on each metric. The price will be applied on the hourly aggregate. Also, it’s common for metrics to be in bytes and not Megabyte or Gigabyte. Keep that in mind when setting the price. For example, if you want to have a price of 1.0$ per GB you will need to set the price to 1.0/10243 = 0.000000000931$ per byte. You don’t have to do this conversion: a rule can carry a price per GB, per GB-hour or per unit, optionally times a multiplier such as 3 replicas, and KCC derives the per-byte factor (see Setting a pricing rule). The pricing rules tab shows every rule as a price, e.g. $0.00012603 per GB-hour × 3 replicas.
Pricing rules are stored in kafka in a compacted topic. The key should be the metric name.
A metric without a pricing rule gets no cost, so costs from pricing rules leave it out. The pricing rules tab lists the collected metrics that have no rule.
Nothing re-creates rules on its own: restarts and redeployments keep them. The sample loader (config-sample) saves every rule in its file when you run it, so running it again overwrites rules with the same metric name, including ones changed or deleted in the UI.
1.4.1. Listing pricing rules
From the UI
Simply go to the pricing rules tab of the UI. You should see the metric name and its cost.
Using Graphql
query getAllRules {
pricingRules {
creationTime
metricName
baseCost
costFactor
}
}
1.4.2. Setting a pricing rule
From the UI
Sign in, go to the pricing rules tab and click Add pricing rule, or the edit button of an existing rule. A metric without a rule can also be priced straight from the list of metrics without a pricing rule. Saving a rule for a metric that already has one replaces it.
Using Graphql
mutation saveRule {
savePricingRule(
request: {
metricName: "confluent_kafka_server_retained_bytes"
baseCost: 0
price: 0.00012603
priceUnit: GB_HOUR
multiplier: 3
multiplierLabel: "replicas"
}
) {
metricName
costFactor
price
priceUnit
multiplier
}
}
priceUnit is what the price is per:
-
GB: per GB of the metric’s value, 1 GB = 10243 bytes (e.g. network bytes); -
GB_HOUR: per GB held for one hour, for stored-bytes metrics, whose hourly maximum is the GB-hours a provider bills; -
UNIT: per unit of the metric’s value for one hour (e.g. per partition).
The optional multiplier applies on top of the price, e.g. 3 for storage that Confluent Cloud bills per replica. KCC derives the costFactor it applies from these: price × multiplier ÷ bytes per unit. You can still send a bare costFactor (per unit of the raw value) instead of a price.
1.4.3. Removing a pricing rule
From the UI
Sign in, go to the pricing rules tab and click the delete button of the rule.
Using Graphql
mutation deleteRule {
deletePricingRule(request: {metricName: "whatever"}) {
creationTime
metricName
baseCost
costFactor
}
}
1.5. Context data
Context data are a way to attach a context (attributes basically) to a kafka item (topic, principal, …). Basically define a set of key/values for an item that match a regex. It is possible that one item match multiple regex (and thus multiple context), but in this case you have to be careful to not have conflicting key/values.
A principal rule’s regex is matched against the principal’s ID and, when the metrics carry one, its display name (Confluent Cloud sends principal_name next to principal_id). So instead of listing service account IDs, one naming-convention rule such as ^kcc-demo-()\.app\.(.)$ with tenant: $1 and application: $2 covers every account that follows the convention, including ones created later. Capture groups come from whichever of the two matched.
You can have as much key/values as you want. They will be used to sum up prices in the dashboard. It is therefor important that you have at least one key/value that defined the cost unit or organization unit. For example: organzation_unit=department1.
The context data are stored in kafka in a compacted topic. The key is free for the user to choose.
1.5.1. Listing existing context data
From the UI
Go to the tab Context Data in the UI. You should see all the context with their validity time, type, regex and context key/values.
If you are signed in, you can use the context tester via the button test context data to check what context key/values currently apply to your topic or principal.
Using Graphql
query getContextData {
contextData {
id
creationTime
validFrom
validUntil
entityType
regex
context {
key
value
}
}
}
1.5.2. Setting context data
If you want to create a new context, you can omit the id if you want. If no id is set, the API will generate one for you using a UUID. If you use an id that is not yet in the system, this means you’re creating a new context item.
From the UI
In the Context Data tab via the button Add context data. You need to be signed in for the button to be visible.
Using Graphql
mutation saveContextData {
saveContextData(
request: {id: "323b603d-5b5f-48d2-84fc-4e784e942289", entityType: TOPIC, regex: ".*collaboration", context: [{key: "app", value: "agoora"}, {key: "cost-unit", value: "spoud"}, {key: "domain", value: "collaboration"}]}
) {
id
creationTime
entityType
regex
context {
key
value
}
}
}
1.5.3. Removing context data
From the UI
Not available yet.
Using Graphql
mutation deleteContextData {
deleteContextData(request: {id: "323b603d-5b5f-48d2-84fc-4e784e942289"}) {
id
creationTime
entityType
regex
context {
key
value
}
}
}
1.6. Reprocess
Reprocessing applies the context rules and pricing rules as they are now to data you already have: for example after fixing a wrong price, or after adding a rule for a service account or topic that showed up without context. Everything from the chosen start time on is rebuilt from the raw metrics; data before it is left as it is.
While the rebuild runs, data from the start time on is incomplete and fills up again: depending on how much raw data there is, this takes minutes to hours. The kcc_* Prometheus gauges replay the historical values during that time.
What happens:
-
Kafka Cost Control stops its Kafka Streams application.
-
It decides where to start: the requested time, rounded down to the start of its aggregation window. If the raw topics no longer reach back that far (retention), it starts at the first complete window they still hold, so nothing is deleted that can’t be rebuilt. Without a start time, it rebuilds everything the raw topics hold.
-
It deletes the application’s internal Kafka Streams topics (
<application id>-…-repartitionand-changelog). They hold the window state and the stream time; left in place, they made replayed windows count as already closed, and they were dropped. -
It deletes the stored windows (the OLAP table) from the start on, so the rebuilt windows replace them instead of being added next to them.
-
It rewinds the raw topics to the start and the pricing rules to their beginning, wipes its local state and restarts (Kubernetes restarts the pod; you’ll see the restart count go up).
-
After the restart it processes the raw data again from the start with today’s rules.
If deleting the internal topics fails (for example because the application’s Kafka user isn’t allowed to delete them), nothing else is changed and the application restarts with its data as it was; the reason is in the response and the logs.
Requirements and side effects:
-
The application’s Kafka user must be allowed to delete its own internal topics (prefix
<application id>-) and to change its consumer group’s offsets. -
The rebuilt windows are written to the
aggregatedandaggregated-table-friendlytopics again. Consumers of those topics (e.g. a JDBC sink) receive new versions of the same keys and should upsert by key. -
Raw data older than the raw topics' retention can’t be reprocessed.
1.6.1. Using the UI
-
Go to the Others tab.
-
Choose a start time (empty means everything the raw topics hold). The quick buttons above help.
-
Click Reprocess and confirm.
Using GraphQL
mutation reprocess {
reprocess(areYouSure: "yes", startTime: "2024-01-01T00:00:00Z")
}
The answer says where the rebuild starts and how many internal topics and stored rows were deleted.
2. Installation
2.1. Prerequisites
This installation manual assumes that
-
You have a Kafka cluster
-
You have a schema registry
-
You have a Kubernetes cluster
2.2. Topics and AVRO schemas
Kafka cost control uses internal topic to compute pricing. You will have to create those topic before deploying the application. The documentation will show the default names, you can change them but don’t forget to adapt the aggregator configuration.
2.2.1. Reference AVRO schemas
Some schemas will reference EntityType. Please add it to your schema registry and reference it when needed.
2.2.2. Topics
| Topic name | Clean up policy | Key | Value |
|---|---|---|---|
context-data |
compact |
String |
|
pricing-rules |
compact |
String |
|
aggregated |
delete |
||
aggregated-table-friendly |
delete |
||
metrics-raw-telegraf-dev |
delete |
None |
String |
Context data
This topic will contain the additional information you wish to attach to the metrics. SEE TODO for more information. This topic is compacted and it is important that you take care of the key yourself. If you wish to delete a context-data you can set null as payload (and provide the key you want to delete).
Pricing rule
This topic will contain the price of each metric. Be aware that most of the metric will be in bytes. So if you want for example to have a price of 1.0$ per GB you will need to set the price to 1.0/10243 = 0.000976276$ per byte. The key should be the metric name. If you wish to remove a price value, send the payload null with the key you want to delete. See TODO on how to use the API or the UI to set the price.
Aggregated
This topic will contain the enriched data. This is the result topic of the aggregator.
Aggregated table friendly
This is the exact same thing as aggregated except there are no hashmap and other nested field. Everything has be flattened. This topic makes it easy to sink the data into a table database.
Metrics raw telegraf
You can have multiple raw topics. For example one per environment or one per kafka cluster. The topic name is up to you, just don’t forget to configure it properly when you deploy telegraf (see Kubernetes section).
Give some special consideration to the retention.ms setting for the raw metrics topics. For example, if you want to distribute the cost of your monthly bill based on the raw metrics scraped over the course of the
month then it is a good idea to retain the scraped data for more than 30 days. This gives people time to ask questions about their bill and also gives the opportunity to reprocess the metrics with new pricing rules/contexts
if needed.
2.3. Kubernetes
You can find all the deployment files in the deployment folder. This folder use Kustomize to simplify the deployment of multiple instances with some variations.
The kubernetes deployment is in two parts. One part is the kafka cost control software (processing and UI) and the other part is the kafka metric scrapper. You may have multiple kafka metric scrapper deployment (one per kafka cluster for example), but you should need only one kafka cost control deployment.
2.3.1. Kafka metric scraper
This part will be responsible to scrape kafka for relevant metrics. Depending on what metrics you want to provide you will need a user with read access to kafka metric but also kafka admin client. Read permission is enough ! You don’t need a user with write permission.
This documentation will assume that you use the dev/ folder, but you can configure as much Kustomize folders as you want. The dev/ folder is a good starting point if you have confluent cluster running: it runs Telegraf on Confluent’s metrics API plus the kafka-scraper (partition and schema counts per topic).
To split a bill’s partition line by tenant, build your overlay on the confluent-scraper folder: the same Telegraf plus the kafka-scraper, which counts each topic’s partitions (dev/ is built on it). Its scraper needs an API key that may describe topics, in the kafka-cost-control-scraper-secret secret; without a Schema Registry, set CC_SCRAPE_SR_ENABLED=false in a kafka-cost-control-scraper-settings config map.
If you only need Confluent’s metrics, build your overlay on the confluent folder instead: Telegraf alone, with the credentials in the telegraf-secret secret. Its non-secret settings (INPUT_URL, OUTPUT_BROKER, OUTPUT_TOPIC, ENVIRONMENT_NAME) can go in a telegraf-settings config map, which takes precedence over the same keys in the secret, so they can live in git.
|
Note
|
Telegraf (1.38 and later) substitutes environment variables only inside quoted strings, so every setting is a plain value: no quotes, no brackets. INPUT_URL is a single URL; to read several clusters, repeat the parameter in it (…?resource.kafka.id=lkc-a&resource.kafka.id=lkc-b). To download only what costs are built from, add metric parameters to INPUT_URL, e.g. &metric=io.confluent.kafka.server/request_bytes; the sample .env lists the six metrics the Telegraf configs keep (network per principal and per topic, storage, partition count). Upgrading from an older version: rename INPUT_URLS to INPUT_URL and OUTPUT_BROKERS to OUTPUT_BROKER, and drop the brackets and quotes from those values and from OUTPUT_TOPIC.
|
Copy the environment sample file:
cd deployment/kafka-metric-scrapper/dev
cp .env.sample .env
vi .env
Edit the environment file with the correct output topic, endpoints and credentials.
Be sure to edit the namespace in the kustomization.yaml file.
Deploy the dev environment using kubectl
cd /deployment/kafka-metric-scrapper
kubectl apply -k dev
Wait for the deployment to finish and check the output topic for metrics. You should receive new data every minute.
2.3.2. Kafka cost control
For this part we will deploy the kafka stream application that is responsible to enrich the metrics (it stores them in its built-in DuckDB database, see Metric storage) and the UI to browse costs and define prices and contexts.
The base folder contains the manifests. Create your own Kustomize overlay next to it, for example deployment/kafka-cost-control/my-env/kustomization.yaml:
apiVersion: kustomize.config.k8s.io/v1beta1
kind: Kustomization
resources:
- ../base
namespace: kafka-cost-control
# the base leaves the Ingress hosts as TO_BE_DEFINED_BY_OVERRIDE, which Kubernetes rejects
patches:
- target:
kind: Ingress
patch: |-
- op: replace
path: /spec/rules/0/host
value: kcc.example.com
generatorOptions:
disableNameSuffixHash: true
secretGenerator:
- name: kafka-cost-control-secret
envs:
- .env
Put your credentials in a .env file in the same folder:
KAFKA_BOOTSTRAP_SERVERS=cluster.region.confluent.cloud:9092
CLUSTER_API_KEY=kafka-api-key
CLUSTER_API_SECRET=kafka-api-secret
KAFKA_SCHEMA_REGISTRY_URL=https://cluster.region.confluent.cloud
SCHEMA_REGISTRY_API_KEY=schema-registry-api-key
SCHEMA_REGISTRY_API_SECRET=schema-registry-api-secret
# Comma separated list of topics that contain raw metrics from telegraf
CC_TOPICS_RAW_DATA=metrics-raw-telegraf-dev
# Admin password
CC_ADMIN_PASSWORD=admin
Replace kcc.example.com with your own host. It is used by both the UI and the aggregator’s /graphql endpoint.
Deploy the application using kubectl
cd deployment/kafka-cost-control
kubectl apply -k my-env
2.3.3. Helm Chart (Strimzi only)
If your Kafka cluster is managed by the Strimzi operators, you can use the provided Helm chart (under helm/kcc-strimzi) to deploy
Cost Control with all the required components (Telegraf, Aggregator, UI) as well as all the required
users and topics. The Helm chart still assumes that you already have a Schema Registry deployed.
See the relevant README file for more information.
A component that is special to the Strimzi deployment is the Context Operator.
This operator detects KafkaUser resources with a specific annotation and automatically creates
a corresponding Context in Cost Control. This allows you to manage your contexts directly from Kubernetes.
See the README in strimzi-operator for more information.
2.4. Metric storage
The aggregator can store the aggregated metrics in a built-in DuckDB database. Parts of the UI (such as grouping in the cost overview), the AI assistant and the exports below read from it.
A stream of aggregated metrics is also always sent to the aggregated and aggregated-table-friendly topics, so you can still feed an external database of your choice from Kafka if you need to.
2.4.1. DuckDB Integration
The provided Kubernetes manifests already enable it. If you deploy the aggregator yourself, set the following environment variables:
containers:
- name: kafka-cost-control
image: spoud/kafka-cost-control:latest
env:
# enable the DuckDB integration
- name: CC_OLAP_ENABLED
value: "true"
# path in the container where the DuckDB database will be stored
# if this is not set, the data will be stored in-memory (i.e. it will be lost when the container is restarted)
- name: CC_OLAP_DATABASE_URL
value: "jdbc:duckdb:/home/jboss/kafka-stream/duckdb.db"
Once the integration is enabled, you can export collected metrics from the DuckDB database by using curl:
# Get all metrics from the last 30 days in CSV format
curl -H "Accept: text/csv" http://localhost:8083/olap/export > out.csv
# Get all metrics from the last 30 days in JSON format
curl -H "Accept: application/json" http://localhost:8083/olap/export > out.json
# Get all metrics for all of February 2025 in CSV format
curl -H "Accept: text/csv" "http://localhost:8083/olap/export?fromDate=2025-02-01T00:00:00Z&toDate=2025-03-01T00:00:00Z" > out.csv
Note that when using the fromDate and toDate parameters, the times must be specified in UTC in the ISO 8601 format (e.g. 2025-02-01T00:00:00Z).
Otherwise the API will return a 400 Bad Request error.
2.5. Troubleshooting
2.5.1. Kafka cost control is never ready
If the kafka-cost-control pod is never ready there are good chances that it is waiting on a topic before it can start. If you look closely in the log you will see a message like this:
2024-02-06 15:46:34,739 WARN [io.qua.kaf.str.run.KafkaStreamsProducer] (pool-5-thread-1) Waiting for topic(s) to be created: [non-existing-topic]
As soon as you create the missing topic(s), you should be good to go. Look again at the Topics section for more information on how to create a topic.
3. Architecture
3.1. Introduction and Goals
Many organizations have introduced Kafka either on premise or in the cloud in recent years.
Kafka platforms are often used as a shared service for multiple teams.
Having all costs centralized in a single cost center means that there is no incentive to save costs for individual users or projects.
Kafka Cost Control gives organizations transparency into the costs caused by applications and allow to distribute platform costs in a fair way to its users by providing a solution that
-
shows usage statistics per application and organizational unit
-
allows defining rules for platform cost distribution over organizational units or applications
-
works for most organizations, no matter if they use Confluent Cloud, Kubernetes or on-prem installations
3.1.1. Requirements Overview
-
Collection and aggregation of usage metrics and statistics from one or multiple Kafka clusters. Aggregation by time:
-
hourly (for debugging or as a metric to understand costs in near real-time)
-
daily
-
weekly
-
monthly
-
-
Management of associations between client applications, projects and organizational units (OU)
-
automatic recognition of running consumer groups
-
automatic detection of principals/clients
-
creation, modification and deletion of contexts (projects and OUs)
-
interface to hook in custom logic for automatic assignment of clients to projects and OUs
-
manual assignment of auto-detected principals or consumer groups to projects and OUs
-
context can change in time, each item should have a start and end date (optional). This means that an item (ex a topic) can switch ownership at any point in time
-
-
Visualization of usage statistics
-
Costs and usage statistics can be broken down interactively
-
Summary view: total costs for timespan (day, week, month) per OU
-
Detail View OU by category: costs by category (produce, consume, storage) for the selected OU in the selected timespan
-
Detail View OU by application/principal/consumer-group/topic
-
-
Data must be made available in a format that can be used to display it with standard software (e.g. Kibana, Grafana, PowerBI), so that organizations can integrate it into an existing application landscape
-
provisioning of a lightweight default dashboard e.g. as a simple SPA, so that extra tooling is not mandatory to view the cost breakdown
-
Items not yet classified should be easily identifiable, so we know what configuration is missing (for example a topic has no OU yet)
-
-
Management of rules, that describe how costs are calculated (aka pricing rules)
-
Management of rules, that describe how costs are calculated, e.g.
-
fixed rates for available metrics, i.e. CHF 0.15 per consumed GB
-
base charge, i.e. CHF 0.5 per principal per hour
-
rules can be changed at any time, but take effect at a specified start time
-
optional: backtesting of rules using historical data
-
-
Access Control
-
only authorized users can modify rules, OUs and projects
-
unauthenticated users should be able to see statistics
-
-
Observability
-
expose metrics so that the cost control app can be monitored
-
proper logging
-
-
Export of end-of-month reports as CSV or Excel for further manual processing
-
Ability to reprocess raw data in case a mistake was made. For example we see at the end of the month that an item was wrongly attributed to an OU. We should be able to correct this and reprocess the data.
3.1.2. Quality Goals
-
Transferability / Extensibility: Kafka Cost Control should be modular, so that company-specific extensions can be added.
A core layer should contain common base functionality. Company specific terms or features should be separated into dedicated modules. -
Maintainability: Reacting to changing requirements and implementing bug fixes should be possible within weeks.
3.1.3. Stakeholders
| Role/Name | Expectations |
|---|---|
Kafka user |
Should be able to see their usage. Should take ownership of resources. |
Management |
Should have an overview of the costs and usage of Kafka. |
3.2. Architecture Constraints
| Constraint | Explanation |
|---|---|
JVM based |
use common language at SPOUD and many clients to make sure many can contribute |
Hosting On-Site (not SaaS only) |
Companies may not want to expose usage data to a SaaS provider |
3.3. System Scope and Context
Kafka Cost Control is a standalone application that needs to integrate into an existing IT landscape.
3.4. Solution Strategy
3.4.1. Used Technologies
| Technology | Reason |
|---|---|
Telegraf |
|
Kafka |
for storing metrics, context info and pricing rules, reduces number of solution dependencies |
Kafka Streams |
for enriching metrics and storing pricing + context data into KTables |
DataStore |
A datastore, e.g. a SQL DB, will be used for the time based aggregations (e.g. end of month reporting). Avoids complex calendar logic in Kafka Streams. |
3.4.2. Time based aggregations & scraping intervals
-
MetricsScraper should ingest metrics with an interval of 1 minute for confluent cloud metrics. Other data sources can have longer intervals.
-
MetricsProcessor aggregates metrics with short time windows of 60 minutes
-
variable cost is usually defined as cost unit/minute
-
The window value is the accumulated cost for one hour (interpolation may be needed when data points are missing)
-
this allows some tolerance for gaps in metrics and varying ingestion intervals
-
3.5. Building Block View
3.5.1. Whitebox Overall System
| Building block | Description |
|---|---|
PricingRules |
Stores rules for turning usage information into costs |
ContextProvider |
Manages contextual information that can be used to enrich metrics with company-specific information. E.g. relations between clientIds, applications, projects, cost centers, … |
MetricProcessor |
|
MetricsScraper |
|
3.5.2. MetricsScraper
Confluent Cloud
Confluent exposes many metrics in prometheus format. These will be scraped with telegraf. Some information are missing from the prometheus export endpoint and need to be fetched with custom queries/requests. This is done with a java application which exposes them as prometheus endpoint. Docs: https://docs.confluent.io/cloud/current/monitoring/metrics-api.html
Additional metric |
endpoint/query |
Partition count of a topic |
|
registered schemas for a topic |
http requests to schema registry needed.
1. |
3.5.3. PricingRules
3.5.4. ContextProvider
Context format
-
metrics are defined in the core
-
a metric belongs to at least one of the dimensions
-
topic
-
consumer group
-
principal
-
-
a context object can be attached to existing dimensions as a AVRO key-value pair to provide the needed flexibility
{
"creationTime": "2024-01-01T00:00:00Z",
"validFrom": "2024-01-01T00:00:00Z",
"validUntil": null,
"entityType": "TOPIC",
"regex": "car-claims",
"context": {
"project": "claims-processing",
"organization_unit": "non-life-insurance",
"sap_psp_element": "1234.234.abc"
}
}
{
"creationTime": "2024-01-01T00:00:00Z",
"validFrom": "2024-01-01T00:00:00Z",
"validUntil": null,
"entityType": "TOPIC",
"regex": "^([a-z0-9-]+)\\.([a-z0-9-]+)\\.([a-z0-9-]+)-.*$",
"context": {
"tenant": "$1",
"app_id": "$2",
"component_id": "$3"
}
}
If naming conventions are very clear they could also be provided as a file / configuration.
{
"creationTime": "2024-01-01T00:00:00Z",
"validFrom": "2024-01-01T00:00:00Z",
"validUntil": null,
"entityType": "PRINCIPAL",
"regex": "u-4j9my2",
"context": {
"project": "claims-processing",
"organization_unit": "non-life-insurance",
"sap_psp_element": "1234.234.abc"
}
}
- INFO
-
Context objects will be started as AVRO messages. We use JSON as a representation here for simplicity.
Context Lookup
State stores in Kafka Streams will be used to construct lookup tables for the context.
The key is a string and is a free value that can be set by the user. If no key is provided the API should create random unique key. The topic is compacted, meaning if we want to delete an item we can send a null payload with its key.
| Key | Value |
|---|---|
<type>_<cluster-id>_<principal_id> |
<context-object> |
PRINCIPAL_lx1dfsg_u-4j9my2_2024-01-01 |
{…, "regex": "u-4j9my2","context": {…}} |
b0bd9c9a-08e6-46c7-9f71-9eafe370da6c |
<context-object> |
Once the table has been loaded, aggregated metrics can be enriched with a KTable - Streams join.
3.6. Runtime View
3.6.1. Metrics Ingestion from Confluent Cloud
Process to gather and aggregate metrics from Confluent Cloud.
The Confluent Metrics Scraper calls the endpoint
api.telemetry.confluent.cloud/v2/metrics/cloud/export?resource.kafka.id={CLUSTER-ID}
with Basic Auth in an interval of 1 Minute to obtain all metrics in Prometheus format.
Telegraf is used to poll data using Confluent prometheus endpoint.
3.6.2. Metrics using Kafka Admin API
Some information can be gathered from the Kafka Admin API. We will develop a simple application that connect to the Kafka Admin API and expose metrics as prometheus endpoint. We can then reuse Telegraf to publish those metrics to kafka.
3.6.3. Other sources of metrics
Anyone can publish to the raw metrics topic. The metrics should follow the telegraf format. Recommendation: use one topic per source of metrics. The MetricEnricher application will anyway consume multiple raw metric topics.
3.6.4. Metrics Enrichment
-
Metrics are consumed from all the raw data topics.
-
Metrics are aggregated by the MetricsProcessor. Here we:
-
aggregate by hours
-
attach context
-
attach pricing rule
-
-
The aggregates are stored in the
aggregated-metricstopic. -
The aggregated metrics are stored into the query database.
The storage procedure into the query database must be idempotent in order to reprocess the enrichment in case of reprocessing.
Enrichment for topics
{
"fields": {
"gauge": 40920
},
"name": "confluent_kafka_server_sent_bytes",
"tags": {
"env": "sdm",
"host": "confluent.cloud",
"kafka_id": "lkc-abc123",
"topic": "agoora-state-global",
"url": "https://api.telemetry.confluent.cloud/v2/metrics/cloud/export?resource.kafka.id=lkc-abc123"
},
"timestamp": 1704805140
}
Enrichment for principals
{
"fields": {
"gauge": 0
},
"name": "confluent_kafka_server_request_bytes",
"tags": {
"env": "sdm",
"host": "confluent.cloud",
"kafka_id": "lkc-abc123",
"principal_id": "u-4j9my2",
"type": "ApiVersions",
"url": "https://api.telemetry.confluent.cloud/v2/metrics/cloud/export?resource.kafka.id=lkc-abc123"
},
"timestamp": 1704805200
}
3.6.5. Metrics Grouping
-
confluent_kafka_server_request_bytes by kafka_id (Cluster) and principal_id (User) for the type Produce as sum stored in produced_bytes
-
confluent_kafka_server_response_bytes by kafka_id (Cluster) and principal_id (User) for the type Fetch as sum stored in fetched_bytes
-
confluent_kafka_server_retained_bytes by kafka_id (Cluster) and topic as min and max stored in retained_bytes_min and retained_bytes_max
-
confluent_kafka_server_consumer_lag_offsets by kafka_id (Cluster) and topic as list of consumer_group_id stored in consumergroups
maybe more are possible.
3.7. Deployment View
3.8. Risks and Technical Debts
-
Difficulty to get context data
-
Will the customer be willing to make the effort to provide the necessary data?
-
-
Difficulty to put a set price on each kafka item
-
How to integrate general cost like operation, etc. (not linked to a particular kafka item)
-
Difficulty of integration with companies cost dashboard
3.9. Glossary
| Term | Definition |
|---|---|
OU |
Organization Unit |