Skip to content

feat(data-warehouse): implement the kafka import source - #108391

Open
Gilbert09 wants to merge 6 commits into
tom/dwh-source-cursorfrom
tom/dwh-kafka-source
Open

Gilbert09 wants to merge 6 commits into
tom/dwh-source-cursorfrom
tom/dwh-kafka-source

Conversation

@Gilbert09

Copy link
Copy Markdown
Member

Problem

Changes

  • Users can connect a Kafka cluster and pick topics. Each topic becomes a table. The source ships as alpha.
  • Authentication options: SASL/PLAIN (Confluent Cloud API keys), SASL/SCRAM-SHA-256, SASL/SCRAM-SHA-512, or none. TLS takes an optional CA certificate.
  • Each sync reads every partition from the stored offsets up to the end offsets captured when it starts. It then stages the next offsets as a KafkaCursor.
  • A partition added since the last sync starts at its first retained message. Starting it at the end would skip what arrived before the sync.
  • A cursor that retention passed, or one beyond the end of a recreated topic, restarts that partition and logs a warning.
  • JSON object fields become columns. Other JSON values go to value, and unparseable messages to _kafka_raw_value, so one bad message does not fail the sync.
  • Every row carries its topic, partition, offset, timestamp, key, headers and a tombstone flag.
  • (partition, offset) is the primary key, so an incremental sync can re-read a window without duplicating rows.
  • _kafka_timestamp is the offered incremental field, because the sync-method picker requires one. The offsets alone decide where a sync starts.
  • The source never joins or commits to a consumer group, so customer consumer groups are untouched.

Warning

The Kafka client connects to every broker the cluster advertises, not only the bootstrap servers. The source checks both against the private-host policy before it consumes. A rejection raises the same HostNotAllowedError the database sources raise, so it stops the sync and stays out of error tracking.

  • Not supported yet: Avro or Protobuf through a Schema Registry, because the client's registry module needs authlib, which is not a dependency.
  • Not supported yet: SSH tunnels, because advertised broker addresses bypass a tunnel.
  • The public docs page is a follow-up in posthog.com, so docsUrl is not set.
  • No frontend code changes. The form renders from the field definitions, like every other source. It was not rendered in a browser.

How did you test this code?

  • Start-offset tests catch a partition that skips data after retention, after a topic recreation, or when it is new.
  • Read-loop tests catch reading past the captured end offsets, and waiting on a partition that ends in transaction markers instead of finishing it at its position.
  • They also catch a stalled partition whose cursor does not stop at its position. Two of these guards were checked by breaking the code and watching the test fail.
  • fetch_cluster tests catch a private advertised broker behind a public bootstrap host, and an authentication failure reported as unreachable.
  • Run locally against the dev stack's Kafka broker, with a 3-partition topic holding JSON, a tombstone and a non-JSON message:
    • The first sync read all 52 messages.
    • The second read only the 12 new messages.
    • A third read none and left the cursor unchanged.
    • A full refresh read all 64 again.
  • Not checked: SASL or TLS against a managed cluster such as Confluent Cloud or MSK, transactional producers on a real broker, and a sync end to end through the warehouse pipeline.

👉 Stay up-to-date with PostHog coding conventions for a smoother review.

Release status

  • No feature flag controls this change
  • This change is behind a feature flag and is not available to users
  • This change makes a previously flagged feature available to everyone

Automatic notifications

  • Publish to changelog?

Docs update

The posthog.com source page is a follow-up.

@Gilbert09 Gilbert09 self-assigned this Sep 29, 2026
@Gilbert09 Gilbert09 added the stamphog Request AI approval (no full review) label Sep 29, 2026 — with Talyn App
@Gilbert09
Gilbert09 added this pull request to stack #108392 September 29, 2026 14:24
@github-actions

Copy link
Copy Markdown
Contributor

Hey @Gilbert09! 👋

It looks like your git author email on this PR isn't your @posthog.com address (owerstom@gmail.com). Since you're on the PostHog team, it's worth pointing your local git author email at your @posthog.com address. Why it matters:

  • Consistent work identity in git history — internal tooling that attributes commits to team members keys off your @posthog.com address.
  • Keeps team contributions easy to tell apart from external community ones when scanning history.

You can fix it for this repo with:

git config user.email "you@posthog.com"

Or set it globally with git config --global user.email "you@posthog.com". No need to redo this PR — just a nudge for next time. 🙂

@github-actions

github-actions Bot commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

🤖 CI report

⚠️ Trunk lane — backend Python lane

This PR is assigned to the backend Python lane. It runs backend Python tests and may merge in parallel with PRs in other lanes.

✅ Duplication (Python) — clean

New Python code duplication introduced by this branch. Fails at 70+ tokens in app code, or 150+ tokens when both copies live in test files. Advisory while the gate proves itself: extract a shared helper instead of copying.

✅ Duplication (TypeScript) — clean

New TypeScript code duplication introduced by this branch. Fails at 70+ tokens in app code, or 150+ tokens when both copies live in test files. Advisory while the gate proves itself: extract a shared helper instead of copying.

✅ Comment density — 3% of added code lines are comments (32 of 1188)

This section warns when comments are more than 3% of the code lines a PR adds, and alerts above 6%. Before agent-assisted PRs, the typical share was about 2%. Only full-line comments count. Docstrings, generated files, snapshots, migrations, and workflow files are left out.

Comments that restate the code, record how the change came about, or narrate the next line add noise for the next reader. Keep the comments that explain a reason the code cannot show, and remove the rest. See .agents/skills/writing-code-comments/SKILL.md for the house rules.

Files with the most added comment lines:

File Comment lines Added lines
products/warehouse_sources/backend/temporal/data_imports/sources/kafka/kafka.py 14 442
products/warehouse_sources/backend/temporal/data_imports/sources/kafka/settings.py 10 40
products/warehouse_sources/backend/temporal/data_imports/sources/kafka/tests/test_kafka.py 5 381
products/warehouse_sources/backend/presentation/views/external_data_source/helpers.py 2 3
products/warehouse_sources/backend/temporal/data_imports/sources/kafka/source.py 1 191

This check does not block merging. It updates on every push and clears when the share drops.

⚠️ Backend coverage — 97.0% of changed backend lines covered — 23 uncovered

🧪 Backend test coverage

Patch coverage — changed backend lines (products + core): ███████████████████░ 97.0% (839 / 862)

File Patch Uncovered changed lines
products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v2/pipeline.py 75.0% 94
products/warehouse_sources/backend/temporal/data_imports/sources/common/typings.py 75.0% 37, 40
products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/pipeline.py 80.0% 128
products/warehouse_sources/backend/temporal/data_imports/sources/kafka/source.py 93.8% 73, 214
products/warehouse_sources/backend/temporal/data_imports/sources/kafka/kafka.py 94.1% 131, 271, 281–282, 373, 380–381, 421, 447–448, 459, 474–476
products/warehouse_sources/backend/temporal/data_imports/sources/common/cursor.py 97.0% 79, 99
products/warehouse_sources/backend/temporal/data_imports/sources/kafka/tests/test_kafka.py 99.6% 187

🤖 Agents: add a test covering the lines above, or note why under "How did you test this code?". Machine-readable gap list: the patch-coverage artifact on this run (gh run download 36598023044 -n patch-coverage), or the coverage-data block at the end of this comment.

Per-product line coverage (touched products)
Product Coverage Lines
experiments ██████░░░░░░░░░░░░░░ 27.8% 8,680 / 31,269
demo ███████████░░░░░░░░░ 52.9% 1,413 / 2,673
batch_exports ████████████░░░░░░░░ 59.0% 15,643 / 26,502
cdp ██████████████████░░ 88.2% 4,545 / 5,155
mcp_analytics ██████████████████░░ 89.2% 5,038 / 5,651
product_tours ██████████████████░░ 89.3% 1,331 / 1,491
dashboards ██████████████████░░ 89.5% 6,904 / 7,714
data_warehouse ██████████████████░░ 90.0% 14,027 / 15,589
notebooks ██████████████████░░ 90.2% 15,287 / 16,945
signals ██████████████████░░ 90.2% 57,818 / 64,071
cohorts ██████████████████░░ 90.4% 8,420 / 9,316
streamlit_apps ██████████████████░░ 90.6% 2,623 / 2,895
managed_warehouse ██████████████████░░ 91.0% 10,252 / 11,263
tasks ██████████████████░░ 91.1% 75,146 / 82,503
data_modeling ██████████████████░░ 91.4% 10,554 / 11,543
exports ██████████████████░░ 91.6% 9,680 / 10,562
engineering_analytics ██████████████████░░ 91.7% 11,032 / 12,030
business_knowledge ██████████████████░░ 92.2% 7,684 / 8,330
conversations ███████████████████░ 92.5% 28,734 / 31,062
early_access_features ███████████████████░ 92.6% 1,332 / 1,439
stamphog ███████████████████░ 92.8% 8,109 / 8,742
canvas ███████████████████░ 92.8% 6,873 / 7,405
approvals ███████████████████░ 93.0% 3,974 / 4,271
mcp_registry ███████████████████░ 93.1% 1,670 / 1,794
notifications ███████████████████░ 93.2% 1,145 / 1,229
slack_app ███████████████████░ 93.2% 13,733 / 14,735
error_tracking ███████████████████░ 93.2% 16,359 / 17,547
surveys ███████████████████░ 93.3% 6,571 / 7,040
autoresearch ███████████████████░ 93.6% 8,481 / 9,061
context_layer ███████████████████░ 93.9% 3,415 / 3,638
web_analytics ███████████████████░ 94.0% 21,815 / 23,218
alerts ███████████████████░ 94.0% 8,570 / 9,114
billing_alerts ███████████████████░ 94.1% 2,094 / 2,226
mcp_store ███████████████████░ 94.4% 8,940 / 9,472
ai_observability ███████████████████░ 94.5% 24,473 / 25,896
workflows ███████████████████░ 94.6% 15,222 / 16,093
wizard ███████████████████░ 94.7% 6,150 / 6,496
reminders ███████████████████░ 94.8% 760 / 802
review_hog ███████████████████░ 95.0% 11,507 / 12,119
endpoints ███████████████████░ 95.1% 9,206 / 9,681
annotations ███████████████████░ 95.1% 817 / 859
customer_analytics ███████████████████░ 95.2% 25,125 / 26,396
marketing_analytics ███████████████████░ 95.3% 19,564 / 20,528
posthog_ai ███████████████████░ 95.3% 2,488 / 2,610
actions ███████████████████░ 95.5% 756 / 792
logs ███████████████████░ 95.5% 15,399 / 16,130
data_catalog ███████████████████░ 95.5% 4,401 / 4,606
tracing ███████████████████░ 95.6% 3,536 / 3,699
growth ███████████████████░ 95.7% 11,228 / 11,734
messaging ███████████████████░ 95.8% 3,798 / 3,963
skills ███████████████████░ 95.8% 6,972 / 7,274
product_analytics ███████████████████░ 96.0% 28,470 / 29,647
revenue_analytics ███████████████████░ 96.4% 1,876 / 1,946
user_interviews ███████████████████░ 96.5% 2,867 / 2,971
feature_flags ███████████████████░ 96.5% 25,588 / 26,509
warehouse_sources ███████████████████░ 97.2% 460,369 / 473,413
data_quality ████████████████████ 97.6% 7,587 / 7,774
metrics ████████████████████ 98.0% 4,084 / 4,166
analytics_platform ████████████████████ 98.3% 2,784 / 2,833
pulse ████████████████████ 98.5% 2,046 / 2,078
live_debugger ████████████████████ 99.2% 626 / 631

Report-only. Patch coverage = changed backend lines covered vs origin/master. Sorted lowest first.
Known gaps: lines covered only by Temporal tests show as uncovered; core line numbers may drift if master changed the same file.

@pr-assigner-resolver-posthog
pr-assigner-resolver-posthog Bot requested a review from a team September 29, 2026 14:24

@stamphog stamphog Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not approved — escalated to a human reviewer.

Re-add the stamphog label to request another review once you have addressed this.

This adds a new outbound network source whose private-host (SSRF) protection is security-sensitive, and it has no independent review and only moderate author familiarity. The host check resolves each name once and then hands the hostname to librdkafka, which resolves it again. The helper's own docs say that leaves a DNS-rebinding gap.

  • Author wrote 33% of the modified lines and has 111 merged PRs in these paths (familiarity MODERATE).
  • SSRF guard is check-then-connect: resolve_safe_host validates hostnames, but librdkafka is given the original hostnames and re-resolves them. resolve_safe_host's docstring says this defeats the check for short-TTL DNS records, and the returned connect_host is ignored.
  • The admin client and consumer can open background connections to the brokers advertised in cluster metadata before or independently of the advertised-broker check in fetch_cluster. The consumer in kafka_source also re-fetches metadata without re-running the broker check.
  • Security-sensitive territory (SSRF protection, credential handling for SASL passwords, TLS CA input) with no agent or human review on the current head.
Gate mechanics and policy version
Gate Result
prerequisites ✓ all clear
deny-list ✓ no deny categories matched
size ✓ 652L, 4F substantive, 1050L/7F incl. docs/generated/snapshots — within ceiling
tier ✓ T1-agent / T1d-complex (1050L, 7F, single-area, feat)
stamphog 2.3.0 .stamphog/policy.yml @ 67c207d · reviewed head 67c207d

@stamphog stamphog Bot removed the stamphog Request AI approval (no full review) label Sep 29, 2026
@veria-ai

veria-ai Bot commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

PR overview

All previously flagged issues have been addressed. No open security concerns remain on this pull request.

Security review

No open security issues remain on this pull request.

Fixed/addressed: 2 · PR risk: 0/10

Copilot AI balanced review requested due to automatic review settings September 29, 2026 14:40

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

@coderabbitai

coderabbitai Bot commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

Review in Change Stack →

Navigate logical layers of code changes, visualize relationships, and explore their blast radius.

Note

Reviews paused

It looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the reviews.auto_review.auto_pause_after_reviewed_commits setting.

Use the following commands to manage reviews:

  • @coderabbitai resume to resume automatic reviews.
  • @coderabbitai review to trigger a single review.

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review
📝 Walkthrough

Walkthrough

The pull request adds Kafka as a warehouse source. It defines connection, authentication, encryption, and message-format settings. It discovers available topics, validates credentials, and imports topic partitions with cursor-based progress. The importer converts JSON and text messages into rows and stops reads at captured partition offsets. A changed topic ID signals that the destination requires a reset. Tests cover configuration, broker discovery, validation, offset handling, row conversion, and partition reads. The source inventory now lists Kafka as implemented.

Priority: ➖ Normal

Merge Risk: 🟡 Moderate · up to a7dd0

Ordinary Cloud teams can select Kafka but cannot connect or import topics. Align its visibility with the supported rollout before merging; also correct the misleading connection error.

Security Architecture Review

Security architecture risk: 🟡 Moderate · up to a7dd0

A Kafka topic replacement triggers a destination rebuild, but an empty retry may leave the old destination intact if the initial reset fails. Kafka access is also deliberately restricted on Cloud. The reset edge case warrants review before rollout.

Retained concerns

  • Medium · reliability · inferred: If a topic-recreation reset fails before removing the old destination, a V3 retry that reads no rows skips the reset, sends no overwrite batch, and can persist the new topic ID. Later runs would no longer request the reset, leaving stale rows in the destination.
Security review details

Security Blast Radius

  • inferred — The network sink is reachable through configured Kafka connections in self-hosted deployments and allowlisted Cloud teams, not through ordinary non-allowlisted Cloud teams. The identified retry path affects the destination table for the importing schema; cross-tenant exposure is not established.

Security Findings and Attack Paths

  • inferred — No attacker-driven network-control bypass is established for non-allowlisted Cloud teams: their requests fail before Kafka client creation. The conditional reset failure could retain rows from a previous topic incarnation, but no cross-tenant access or verified Security finding is established.

Trust Boundaries and Controls

  • observed — The source checks bootstrap and advertised broker hosts, while a separate early team gate prevents non-allowlisted Cloud callers from reaching librdkafka. Hostname validation alone would not bind a later client lookup to the checked address.

Resilience and Maintainability Implications

  • inferred — Staging protects nonempty loads from advancing the durable cursor before loader completion, but the direct zero-batch cursor commit makes successful reset completion a necessary precondition when a topic identity changes.

Hardening Proposals

  • proposed — Preserve a pending topic-identity reset across retries until destination reset or overwrite is confirmed, including when the replacement topic yields zero batches.
  • proposed — Before extending Kafka access to ordinary Cloud teams, provide a way to validate the addresses actually dialed by the Kafka client; the current hostname preflight does not provide that guarantee.
🚥 Pre-merge checks | ✅ 1
✅ Passed checks (1 passed)
Check name Status Explanation
Description check ✅ Passed The description is complete and standalone. It explains the problem, user-visible changes, supported and unsupported behavior, testing performed, limitations, release status, notifications, and docume…
✨ Finishing Touches 💡 1
🛠️ Fix failing CI checks 💡
  • Commit to this branch
  • Create a new PR
📝 Generate docstrings
  • Commit to this branch
  • Create a new PR
🧪 Generate unit tests (beta)
  • Commit to this branch
  • Create a new PR

Comment @coderabbitai help to get the list of available commands.

@trunk-io

trunk-io Bot commented Sep 29, 2026 •

Copy link
Copy Markdown

Static Badge   Static Badge   Static Badge

Failed Test Failure Summary Logs
test_reserved_and_invalid_payload_fields_are_namespaced The test failed because no keys in the data started with the expected prefix '_kafka_payload_not_valid_'. Logs ↗︎

View Full Report ↗︎ ⋅ Docs

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Note

Quiet mode is enabled, so only the most important comments were posted inline. Other review comments are grouped below.

🟡 Other comments (1)
products/warehouse_sources/backend/temporal/data_imports/sources/kafka/kafka.py-369-375 (1)

369-375: 🗄️ Data Integrity & Integration | 🟡 Minor | ⚡ Quick win

Reject errored topic metadata before staging the cursor.

If list_topics returns the topic with error set and no partitions, watermarks becomes {}. read_partitions then stages {}. The Kafka source uses the default cursor merge, so this replaces a stored non-empty cursor. The zero-batch pipeline persists that staged cursor directly. A later run starts each partition at its low watermark and can reread retained messages.

KafkaException from list_topics or get_watermark_offsets already closes the consumer. A missing topics[topic] entry instead raises KeyError outside that handler and leaks the consumer. Handle both missing and errored metadata:

Suggested fix
-        topic_metadata = consumer.list_topics(topic, timeout=METADATA_TIMEOUT_SECONDS).topics[topic]
+        topic_metadata = consumer.list_topics(topic, timeout=METADATA_TIMEOUT_SECONDS).topics.get(topic)
+        if topic_metadata is None or topic_metadata.error is not None:
+            consumer.close()
+            raise KafkaSourceError(TOPIC_NOT_FOUND_MESSAGE)
         watermarks = {

ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository: PostHog/posthog/.coderabbit.yaml

Review profile: QUIET

Plan: Enterprise

Run ID: 3d2b3412-4dbb-4e92-98af-2225c3cb1632

📥 Commits

Reviewing files that changed from the base of the PR and between 5899e56 and dd137e2.

📒 Files selected for processing (8)
  • products/warehouse_sources/backend/temporal/data_imports/sources/SOURCES.md
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/test/__init__.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/generated_configs/kafka.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/kafka/kafka.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/kafka/settings.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/kafka/source.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/kafka/tests/test_kafka.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/kafka/tests/test_kafka_source.py

Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 2 remain after this review.

@Gilbert09
Gilbert09 requested a review from a team September 29, 2026 14:50

@fuziontech fuziontech left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Automated review generated on behalf of @fuziontech. I reviewed the complete diff against the requested base and ran compile checks; focused pytest/lint commands were unavailable because this checkout has no configured Python 3.14.7 environment or pytest/ruff binaries. No blocking issue found, so approving.

Non-blocking follow-ups:

  • The Kafka host preflight validates bootstrap and advertised broker DNS results but passes the original hostnames to librdkafka, which resolves them again. The current Cloud guard limits this path to allowlisted teams, but pinning the validated addresses (or documenting the trusted exception) would make the SSRF guarantee explicit.
  • Read-time topic ACL failures can surface as a raw KafkaException from read_partitions; mapping authorization errors to the existing actionable non-retryable Kafka error would avoid retry loops when a user can describe a topic but cannot consume it.

Each topic becomes a table. A sync reads every partition from the offsets
the last sync stopped at up to the end offsets captured when it starts, and
stores the next offsets as a KafkaCursor once the rows are durable. New
partitions start at their first retained message. A cursor that retention
passed, or that is past a recreated topic's end, restarts that partition.

Values are JSON (object fields become columns) or text. Each row carries its
topic, partition, offset, timestamp, key, headers and a tombstone flag.
(partition, offset) is the primary key, so an incremental sync can re-read a
window without duplicating rows.

The source checks the bootstrap servers and every broker the cluster
advertises against the private-host policy, because the client connects to
the advertised brokers. It never joins or commits to a consumer group.
Block untrusted PostHog Cloud Kafka connections until librdkafka can pin DNS resolutions, and fix the warehouse-source pytest module collision.
Treat Kafka authentication as a nested credential container so changing brokers requires password re-entry, with API regression coverage.

Copy link
Copy Markdown
Member Author

All CI checks pass on the latest head, every review thread is resolved, and the branch is current with and cleanly mergeable into its base. GitHub still reports the PR as blocked pending the required human review; approval is needed from Team Warehouse Sources.

🦉 via talyn.dev

Address Kafka credential retargeting, transactional isolation, topic recreation, retry idempotency, metadata failures, binary preservation, timestamp bounds, and reserved payload columns.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Note

Quiet mode is enabled, so only the most important comments were posted inline. Other review comments are grouped below.

Caution

Some comments are outside the diff and can’t be posted inline due to GitHub limitations.

⚠️ Outside diff range comments (1)

🟠 Major · Align Kafka visibility with the Cloud guard. · kafka.py:181-186

products/warehouse_sources/backend/temporal/data_imports/sources/kafka/kafka.py:181-186
🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

Align Kafka visibility with the Cloud guard.

KafkaSource is globally registered and marked only as ReleaseStatus.ALPHA; the repository defines alpha sources as exposed to users. No feature flag or other rollout gate limits it to the two allowlisted analytics teams.

A normal Cloud team can therefore select Kafka, but fetch_cluster() rejects it before creating a Kafka client. This makes credential validation fail, prevents topic discovery, and blocks imports because kafka_source() calls fetch_cluster() again.

Add a rollout gate that hides Kafka from non-allowlisted Cloud teams, or implement broker-address pinning before permitting ordinary Cloud teams. Do not remove this guard alone, because the code documents that librdkafka can bypass a separate DNS preflight through broker re-resolution.

🟡 Other comments (1)
products/warehouse_sources/backend/temporal/data_imports/sources/kafka/kafka.py-441-449 (1)

441-449: 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win

Preserve the Kafka error classification for describe_topics.

describe_topics(...).result() can raise KafkaException for failures other than UNKNOWN_TOPIC_OR_PART. The current catch maps every such exception to TOPIC_NOT_FOUND_MESSAGE, even though the same client reports authentication and TLS failures through errors. Preserve the topic-not-found message only for UNKNOWN_TOPIC_OR_PART; otherwise use errors.user_error().

🐛 Suggested fix
         except KafkaException as e:
-            raise KafkaSourceError(TOPIC_NOT_FOUND_MESSAGE) from e
+            if e.args and e.args[0].code() == KafkaError.UNKNOWN_TOPIC_OR_PART:
+                raise KafkaSourceError(TOPIC_NOT_FOUND_MESSAGE) from e
+            raise errors.user_error() from e
🧹 Nitpick comments (1)
products/warehouse_sources/backend/tests/api/test_external_data_source.py (1)

6872-6902: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Consolidate the overlapping Kafka rejection coverage.

The later test protects one additional assertion: the stored password remains unchanged. Move that assertion into test_kafka_host_change_requires_nested_password, then remove the later test. The fixtures contain different optional fields, but both rejection requests exercise the same nested-password requirement.

Suggested consolidation
         source.refresh_from_db()
         assert source.job_inputs["bootstrap_servers"] == "broker.example.com:9092"
+        assert source.job_inputs["authentication"]["password"] == "stored-secret"
         mock_validate_credentials.assert_not_called()

ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository: PostHog/posthog/.coderabbit.yaml

Review profile: QUIET

Plan: Enterprise

Run ID: 26454993-5853-4d8a-b159-bcbf046a31d7

📥 Commits

Reviewing files that changed from the base of the PR and between 298347e and a7dd078.

📒 Files selected for processing (8)
  • products/warehouse_sources/backend/temporal/data_imports/sources/common/typings.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/kafka/kafka.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/kafka/settings.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/kafka/source.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/kafka/tests/test_kafka.py
  • products/warehouse_sources/backend/temporal/data_imports/sources/kafka/tests/test_kafka_source.py
  • products/warehouse_sources/backend/temporal/data_imports/workflow_activities/import_data_sync.py
  • products/warehouse_sources/backend/tests/api/test_external_data_source.py

Included review availability: This review used your included allowance. Your plan provides up to 12 included reviews per hour; 2 remain after this review.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants