This release is a pre-release and may not be stable for production use.
conduit-connector-sdk (Python)
Python SDK for building Conduit source and destination connectors. Connectors built with this SDK run as standalone gRPC subprocess plugins — no changes to Conduit itself are required to run one.
Status: pre-alpha, under active development. No release has shipped yet. The public API (
Source,Destination,Config,Record) is not stable until the v0.19 Phase-1 acceptance criteria indocs/design/20260707-python-connector-sdk.mdare met. Treat everything here as subject to change without a deprecation notice until then.What "GA" means for this repo's first tagged release (v0.20 plan, WS2), stated precisely so it isn't over-read: GA here means exactly two things -- the acceptance suite (
conduit.testing.acceptance) is green in CI, and a first tag is published via PyPI Trusted Publishing. It does not mean production hardening, edge-case parity with the Go SDK, or a 1.0. The value of that tag is narrow and specific: it unblocks the Rust SDK build, which waits on a tagged release existing, not on a merge tomain. Read every claim in this README with that scope in mind.
What this is
- A gRPC client/server implementation of
conduit-connector-protocolv2 (SourcePlugin/DestinationPlugin/SpecifierPlugin), wrapped in an idiomatic Python API:async/awaitconnector methods (with sync-method auto-dispatch to a thread pool),pydantic-based config with automatic parameter introspection, and abytes | dictOpenCDC record model instead of Go'sDatainterface. - The Python analog of
conduit-connector-sdk(Go). Behavioral parity is a goal; API-shape parity is not — see the design doc's "Alternatives considered" section for why.
Batching (sdk.batch.size / sdk.batch.delay)
Every connector gets two SDK-injected config parameters for free -- authors
never declare them on their own BaseConfig -- matching the Go SDK's
DestinationWithBatch/SourceWithBatch middleware exactly:
sdk.batch.size(int, default0): maximum records per batch before an early flush.sdk.batch.delay(Go-duration string, default0s): maximum time an incomplete batch waits before flushing.
Activation threshold, matched exactly from the Go SDK: batching is only
active when size > 1 or delay > 0. A size of 0 or 1 with no
delay is passthrough -- no accumulation, no background task, one wire
message in for one wire message/record out (Go's destination.go:182/
source_middleware.go:709 condition, reproduced verbatim, not rounded to a
more "intuitive" size >= 1).
On the source side, this SDK's Source ABC has no read_batch/ReadN
override point yet (deferred past this phase), so batching buffers
individual read() calls -- Go's own fallback behavior when a connector
doesn't implement ReadN, not its optimized path. This is a real, current
capability gap versus Go, documented (not hidden) in
conduit/_batch.py's module docstring.
On the destination side, incoming Run batches are flattened and
re-grouped by size/delay before write() is called -- so write() may now
receive a batch spanning several incoming wire messages. A buffered-but-
below-threshold remainder is always flushed (written and acked) when the
stream ends or Stop() drains the connector -- batching never drops
records, per invariant 3 (at-least-once is the floor).
Avro schema support
conduit.schema.AvroSchema (optional -- install with
conduit-connector-sdk[avro]) encodes/decodes plain, headerless Avro
binary -- the exact wire format conduit-commons' schema/avro package
produces via github.com/hamba/avro/v2 (not the Confluent Schema
Registry wire format; there's no magic-byte/schema-ID header on either
side). This is a cross-language compatibility claim, proven against golden
bytes an actual Go program produced with hamba/avro/v2 directly, not a
Python-only round-trip -- see tests/test_schema_avro.py and
tests/testdata/avro_golden.json.
This SDK does not ship a schema registry client (no SchemaService
gRPC stubs are generated) -- an author supplies schema text themselves and
is responsible for keeping it consistent with the
opencdc.{key,payload}.schema.{subject,version} record metadata (see
conduit.record.Metadata.set_payload_schema/get_payload_schema).
What this is not (yet)
See the design doc's Phase 2/3 breakdown for what's still deliberately
deferred: a read_batch/ReadN override point for source connectors (see
"Batching" above for the current fallback-only behavior), a schema registry
client, the acceptance-test harness's full test corpus beyond the current
categories, and any performance claim versus the Go SDK (none ships without
a committed benchi result, per the org's CLAUDE.md).
Delivery semantics
What this SDK guarantees, and — just as important — what it does not:
Guaranteed:
- At-least-once. A source record is never acknowledged to the
connector's own
ack()hook until Conduit'sRunstream sends its position back viaack_positions(never speculatively when the record is merely produced) -- seeconduit.source._SourceServicer._consume_acks. - No silent partial-batch acking. A destination write failure never
assumes an unmentioned record succeeded --
write()either returns cleanly (full-batch success), raisesBatchWriteErrorwith an exhaustive, construction-time-validated per-index accounting, or (any other exception) nacks the entire batch. There is no code path that infers "everything not explicitly marked as failed" succeeded — seeconduit.errors.BatchWriteErrorandconduit.destination._DestinationServicer._write_batch. - Batching never drops records. A buffered-but-below-threshold
remainder is always flushed (and, on the destination side, acked) before
a controlled shutdown (stream end,
Stop(), or a SIGTERM-triggered drain) completes. - Graceful shutdown by default. SIGTERM drains an in-flight read/write
loop (including any buffered batch) before
teardown()runs, bounded by a watchdog so a genuinely wedged connector still exits.
Not guaranteed / explicit limits:
- Positions/state are the connector author's responsibility. This SDK
round-trips whatever
bytesaSource.read()/open()implementation returns; it does not itself provide crash-safe position storage. - A hard-cancelled (not drained)
Runcan lose a buffered-but-unflushed batching remainder. The controlled shutdown paths above (stream end,Stop(), SIGTERM-triggered drain) all flush correctly; an externally torn-down RPC (e.g. the framework cancelling the whole call, not a normal drain) is a documented, narrow-window exception -- seeconduit._batch.collect_batches's docstring for exactly why and how this mirrors an already-accepted tradeoff elsewhere in this SDK (an in-flight, never-acked write on a hard stop). - No exactly-once. Like the Go SDK, redelivery on restart is possible; connectors must tolerate at-least-once delivery.
- Structured (
dict) payloads lose precision on integers beyond2**53(viagoogle.protobuf.Struct's double-precision representation) -- silently, by design of the wire format, not a bug in this SDK. Pinned bytests/test_record_codec.py. - No schema registry integration, no schema evolution/compatibility checking -- see "Avro schema support" above.
Requirements
- Python 3.11+
uvfor dependency management (recommended;pip install -e .[dev]also works)
Repo layout
src/conduit/
__init__.py # public API surface
config.py # BaseConfig, Field, to_parameters()
record.py # Record / Change / Operation / Metadata
schema.py # AvroSchema -- optional, `conduit-connector-sdk[avro]`
source.py # Source ABC
destination.py # Destination ABC
_batch.py # sdk.batch.size/delay middleware (source + destination)
serve.py # handshake + gRPC server bootstrap
_handshake.py # magic cookie, protocol negotiation, stdout line
_build.py # `conduit-connector-sdk build` implementation
_cli.py # `conduit-connector-sdk` console-script entry point
_grpc/ # generated protobuf/grpc stubs (buf generate output)
testing/ # acceptance-test harness (acceptance.py, fixtures.py)
examples/http-poll-source/ # worked example connector
docs/design/ # design docs for this repo
tests/ # unit tests
Building a standalone connector artifact
Conduit launches a standalone connector as a subprocess with a clean
environment — no inherited PATH (design doc §1.1.6). A pip install-then-shebang-script connector (#!/usr/bin/env python3, or an
activated venv) cannot launch this way: there's no PATH for env to
search. conduit-connector-sdk build closes that gap:
conduit-connector-sdk build examples/http-poll-source -o http-poll-source.pyz
./http-poll-source.pyz # directly executable — no `python` prefix, no venv activation
This produces one file with an absolute interpreter shebang (resolved
at build time, never looked up via PATH), bundling every third-party
dependency your connector needs — including compiled-extension
dependencies like grpcio/pydantic's pydantic-core, which a plain
zipapp can't load
in-place: the artifact extracts itself to a per-build cache directory on
first run (the same fundamental approach shiv/pex use), then executes
your connector's real entry point from those extracted files. Later
launches of the same build reuse the cache.
Precondition: run build from an environment where your connector's
own dependencies are already installed (however you installed them — pip,
uv, poetry) — it vendors from what's already resolved, not a fresh
pip install. See conduit/_build.py's module docstring for the full
rationale and known limitations.
Contributing
See CONTRIBUTING.md. This SDK sits on Conduit's data path —
read docs/design/20260707-python-connector-sdk.md
before proposing changes to the wire adapter, ack/nack logic, or handshake.
License
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file conduit_connector_sdk-0.1.0.dev1.tar.gz.
File metadata
- Download URL: conduit_connector_sdk-0.1.0.dev1.tar.gz
- Upload date:
- Size: 94.1 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via:
twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
2491e8a946db8358b5e9545b2f8e7aea4174da76798a849f51fc99548b614d10
|
|
| MD5 |
286756d631f9d323045c07ed276882fa
|
|
| BLAKE2b-256 |
d006e4016b2807500a2df0771f5dbbc6681bf1d3e4fa8b4d825702da4720a68a
|
Provenance
The following attestation bundles were made for conduit_connector_sdk-0.1.0.dev1.tar.gz:
Publisher:
release.yml on ConduitIO/conduit-connector-sdk-python
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
conduit_connector_sdk-0.1.0.dev1.tar.gz -
Subject digest:
2491e8a946db8358b5e9545b2f8e7aea4174da76798a849f51fc99548b614d10 - Sigstore transparency entry: 2642071098
- Sigstore integration time:
-
Permalink:
ConduitIO/conduit-connector-sdk-python@c4befc2da1f0e7c8561ae8c0b886567670408462 -
Branch / Tag:
refs/tags/v0.1.0.dev1 - Owner: https://github.com/ConduitIO
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
release.yml@c4befc2da1f0e7c8561ae8c0b886567670408462 -
Trigger Event:
push
-
Statement type:
File details
Details for the file conduit_connector_sdk-0.1.0.dev1-py3-none-any.whl.
File metadata
- Download URL: conduit_connector_sdk-0.1.0.dev1-py3-none-any.whl
- Upload date:
- Size: 111.4 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via:
twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
6927599db6867ebd3cfa0ad3752a09afa38460547729b4c86b15b3710e773fd2
|
|
| MD5 |
cdcf53578ce90a7c63848c7578d51a76
|
|
| BLAKE2b-256 |
969ef7eb8eb5f89efa17f07448db036a10ad9201210f92c18a77b9093d5dbae7
|
Provenance
The following attestation bundles were made for conduit_connector_sdk-0.1.0.dev1-py3-none-any.whl:
Publisher:
release.yml on ConduitIO/conduit-connector-sdk-python
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
conduit_connector_sdk-0.1.0.dev1-py3-none-any.whl -
Subject digest:
6927599db6867ebd3cfa0ad3752a09afa38460547729b4c86b15b3710e773fd2 - Sigstore transparency entry: 2642071178
- Sigstore integration time:
-
Permalink:
ConduitIO/conduit-connector-sdk-python@c4befc2da1f0e7c8561ae8c0b886567670408462 -
Branch / Tag:
refs/tags/v0.1.0.dev1 - Owner: https://github.com/ConduitIO
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
release.yml@c4befc2da1f0e7c8561ae8c0b886567670408462 -
Trigger Event:
push
-
Statement type: