Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 0 additions & 7 deletions nodejs/src/cdp/outputs/producer-registry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,6 @@ import { KafkaProducerRegistryBuilder } from '~/common/outputs/kafka-producer-re
import {
WAREHOUSE_PRODUCER,
WAREHOUSE_PRODUCER_CONFIG_MAP,
WARPSTREAM_CALCULATED_EVENTS_PRODUCER,
WARPSTREAM_CALCULATED_EVENTS_PRODUCER_CONFIG_MAP,
WARPSTREAM_CYCLOTRON_PRODUCER,
WARPSTREAM_CYCLOTRON_PRODUCER_CONFIG_MAP,
WARPSTREAM_INGESTION_PRODUCER,
Expand All @@ -16,10 +14,6 @@ import {
*
* - `WARPSTREAM_INGESTION_PRODUCER` — hog function monitoring (app metrics +
* log entries). Targets the warpstream-ingestion cluster.
* - `WARPSTREAM_CALCULATED_EVENTS_PRODUCER` — dedicated calculated-events
* cluster. No CDP output routes to it now that the precalculated-filters
* consumer is gone; kept registered because charts still supply its env vars
* and registration is lazy. Reuse it for the next calculated-events output.
* - `WARPSTREAM_CYCLOTRON_PRODUCER` — Cyclotron Warpstream cluster used for
* batch hogflow request enqueue. Distinct env-var prefix from the legacy
* `KAFKA_CDP_PRODUCER_*` so output routing is decoupled from the cyclotron
Expand All @@ -29,7 +23,6 @@ import {
export function createCdpProducerRegistry(kafkaClientRack: string | undefined) {
return new KafkaProducerRegistryBuilder(kafkaClientRack)
.register(WARPSTREAM_INGESTION_PRODUCER, WARPSTREAM_INGESTION_PRODUCER_CONFIG_MAP)
.register(WARPSTREAM_CALCULATED_EVENTS_PRODUCER, WARPSTREAM_CALCULATED_EVENTS_PRODUCER_CONFIG_MAP)
.register(WARPSTREAM_CYCLOTRON_PRODUCER, WARPSTREAM_CYCLOTRON_PRODUCER_CONFIG_MAP)
.register(WAREHOUSE_PRODUCER, WAREHOUSE_PRODUCER_CONFIG_MAP)
}
80 changes: 1 addition & 79 deletions nodejs/src/cdp/outputs/producers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,14 +15,6 @@ export type WarpstreamIngestionProducer = typeof WARPSTREAM_INGESTION_PRODUCER
export const WAREHOUSE_PRODUCER = 'WAREHOUSE_PRODUCER' as const
export type WarehouseProducer = typeof WAREHOUSE_PRODUCER

/**
* Targets the dedicated Warpstream calculated-events cluster
* (`KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_*`). Currently unused — see
* `createCdpProducerRegistry` for why it stays registered.
*/
export const WARPSTREAM_CALCULATED_EVENTS_PRODUCER = 'WARPSTREAM_CALCULATED_EVENTS_PRODUCER' as const
export type WarpstreamCalculatedEventsProducer = typeof WARPSTREAM_CALCULATED_EVENTS_PRODUCER

/**
* Targets the Warpstream Cyclotron cluster for outputs (batch hogflow requests).
* Distinct from the legacy `KAFKA_CDP_PRODUCER_*` env vars used by the cyclotron
Expand All @@ -33,11 +25,7 @@ export const WARPSTREAM_CYCLOTRON_PRODUCER = 'WARPSTREAM_CYCLOTRON_PRODUCER' as
export type WarpstreamCyclotronProducer = typeof WARPSTREAM_CYCLOTRON_PRODUCER

/** Producer names registered by the CDP deployments. */
export type CdpProducerName =
| WarpstreamIngestionProducer
| WarehouseProducer
| WarpstreamCalculatedEventsProducer
| WarpstreamCyclotronProducer
export type CdpProducerName = WarpstreamIngestionProducer | WarehouseProducer | WarpstreamCyclotronProducer

/** Cluster-named env var prefix mirroring the producer it configures. */
export const WARPSTREAM_INGESTION_PRODUCER_CONFIG_MAP = {
Expand Down Expand Up @@ -190,72 +178,6 @@ export function getDefaultKafkaWarehouseProducerEnvConfig(): KafkaWarehouseProdu
}
}

/**
* Warpstream calculated-events producer. Targets the dedicated cluster for
* `clickhouse_prefiltered_events` + `clickhouse_precalculated_person_properties`.
*/
export const WARPSTREAM_CALCULATED_EVENTS_PRODUCER_CONFIG_MAP = {
'client.id': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_CLIENT_ID',
'metadata.broker.list': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_METADATA_BROKER_LIST',
'security.protocol': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_SECURITY_PROTOCOL',
'compression.codec': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_COMPRESSION_CODEC',
'linger.ms': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_LINGER_MS',
'batch.size': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_BATCH_SIZE',
'queue.buffering.max.messages': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_QUEUE_BUFFERING_MAX_MESSAGES',
'queue.buffering.max.kbytes': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_QUEUE_BUFFERING_MAX_KBYTES',
'enable.idempotence': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_ENABLE_IDEMPOTENCE',
'message.max.bytes': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_MESSAGE_MAX_BYTES',
'batch.num.messages': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_BATCH_NUM_MESSAGES',
'sticky.partitioning.linger.ms': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_STICKY_PARTITIONING_LINGER_MS',
'topic.metadata.refresh.interval.ms':
'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_TOPIC_METADATA_REFRESH_INTERVAL_MS',
'metadata.max.age.ms': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_METADATA_MAX_AGE_MS',
'message.send.max.retries': 'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_RETRIES',
'max.in.flight.requests.per.connection':
'KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION',
} as const satisfies Partial<Record<AllowedConfigKey, string>>

/** Typed env vars referenced by `WARPSTREAM_CALCULATED_EVENTS_PRODUCER_CONFIG_MAP`. */
export type KafkaWarpstreamCalculatedEventsProducerEnvConfig = {
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_CLIENT_ID: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_METADATA_BROKER_LIST: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_SECURITY_PROTOCOL: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_COMPRESSION_CODEC: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_LINGER_MS: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_BATCH_SIZE: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_QUEUE_BUFFERING_MAX_MESSAGES: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_QUEUE_BUFFERING_MAX_KBYTES: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_ENABLE_IDEMPOTENCE: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_MESSAGE_MAX_BYTES: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_BATCH_NUM_MESSAGES: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_STICKY_PARTITIONING_LINGER_MS: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_TOPIC_METADATA_REFRESH_INTERVAL_MS: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_METADATA_MAX_AGE_MS: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_RETRIES: string
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION: string
}

export function getDefaultKafkaWarpstreamCalculatedEventsProducerEnvConfig(): KafkaWarpstreamCalculatedEventsProducerEnvConfig {
return {
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_CLIENT_ID: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_METADATA_BROKER_LIST: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_SECURITY_PROTOCOL: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_COMPRESSION_CODEC: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_LINGER_MS: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_BATCH_SIZE: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_QUEUE_BUFFERING_MAX_MESSAGES: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_QUEUE_BUFFERING_MAX_KBYTES: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_ENABLE_IDEMPOTENCE: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_MESSAGE_MAX_BYTES: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_BATCH_NUM_MESSAGES: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_STICKY_PARTITIONING_LINGER_MS: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_TOPIC_METADATA_REFRESH_INTERVAL_MS: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_METADATA_MAX_AGE_MS: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_RETRIES: '',
KAFKA_WARPSTREAM_CALCULATED_EVENTS_PRODUCER_MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION: '',
}
}

/**
* Warpstream Cyclotron output producer. Distinct env-var prefix from the legacy
* `KAFKA_CDP_PRODUCER_*` so output routing is decoupled from the cyclotron job
Expand Down
2 changes: 0 additions & 2 deletions nodejs/src/common/config/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@ import { getDefaultAIObservabilityConfig } from '~/ai-observability/config'
import { getDefaultCdpConfig } from '~/cdp/config'
import {
getDefaultKafkaWarehouseProducerEnvConfig,
getDefaultKafkaWarpstreamCalculatedEventsProducerEnvConfig,
getDefaultKafkaWarpstreamCyclotronProducerEnvConfig,
getDefaultKafkaWarpstreamIngestionProducerEnvConfig,
} from '~/cdp/outputs/producers'
Expand All @@ -25,7 +24,6 @@ export function getDefaultConfig(): PluginsServerConfig {
...getDefaultLogsIngestionConsumerConfig(),
...getDefaultTracesIngestionConsumerConfig(),
...getDefaultKafkaWarpstreamIngestionProducerEnvConfig(),
...getDefaultKafkaWarpstreamCalculatedEventsProducerEnvConfig(),
...getDefaultKafkaWarpstreamCyclotronProducerEnvConfig(),
...getDefaultKafkaWarehouseProducerEnvConfig(),
}
Expand Down
2 changes: 0 additions & 2 deletions nodejs/src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@ import type { AIObservabilityConfig } from './ai-observability/config'
import type { CdpConfig } from './cdp/config'
import type {
KafkaWarehouseProducerEnvConfig,
KafkaWarpstreamCalculatedEventsProducerEnvConfig,
KafkaWarpstreamCyclotronProducerEnvConfig,
KafkaWarpstreamIngestionProducerEnvConfig,
} from './cdp/outputs/producers'
Expand Down Expand Up @@ -106,7 +105,6 @@ export interface PluginsServerConfig
TracesIngestionConsumerConfig,
// Producer envs needed by the CDP producer registry the legacy big server builds.
KafkaWarpstreamIngestionProducerEnvConfig,
KafkaWarpstreamCalculatedEventsProducerEnvConfig,
KafkaWarpstreamCyclotronProducerEnvConfig,
KafkaWarehouseProducerEnvConfig {}

Expand Down
2 changes: 0 additions & 2 deletions nodejs/tests/helpers/cdp.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,6 @@ import { CdpConsumerBaseDeps } from '../../src/cdp/consumers/cdp-base.consumer'
import {
CdpProducerName,
WAREHOUSE_PRODUCER,
WARPSTREAM_CALCULATED_EVENTS_PRODUCER,
WARPSTREAM_CYCLOTRON_PRODUCER,
WARPSTREAM_INGESTION_PRODUCER,
} from '../../src/cdp/outputs/producers'
Expand All @@ -28,7 +27,6 @@ function buildTestCdpProducerRegistry(
): KafkaProducerRegistry<CdpProducerName> {
return new KafkaProducerRegistry<CdpProducerName>({
[WARPSTREAM_INGESTION_PRODUCER]: kafkaProducer,
[WARPSTREAM_CALCULATED_EVENTS_PRODUCER]: kafkaProducer,
[WARPSTREAM_CYCLOTRON_PRODUCER]: kafkaProducer,
[WAREHOUSE_PRODUCER]: kafkaProducer,
})
Expand Down
Loading