Skip to main content

confluent-sql

A DB-API v2 compliant Python driver for Confluent Cloud Flink SQL services.

Overview

The confluent-sql library provides a standard DB-API v2 interface for connecting to and executing SQL queries against Confluent Cloud Flink SQL services. This allows you to use familiar database programming patterns with Confluent's streaming SQL capabilities.

Status

This is pre-production code mainly developed as the lower level portion of a dbt adaptor for Confluent Cloud Flink, but is aimed to also be a reasonable standalone dbapi+ driver for python programs to interact with Confluent Flink SQL.

The behavior of snapshot-mode cursors, complying with dbapi semantics, are well stable. The streaming query extensions are more of a work in progress at this time. Feedback and suggestions are welcome!

The driver defaults to snapshot mode for all queries unless streaming mode is explicitly requested.

Prerequisites

  • Confluent Cloud account with Flink environment

  • Existing Flink Database (Confluent Cloud Kafka cluster)

  • API credentials, one of:

    • a "Global" Confluent Cloud API key and secret, passed as global_api_key / global_api_secret, or
    • a Flink Region API key and secret, passed as flink_api_key / flink_api_secret, or
    • a BYOIDC bearer token minted by your own OAuth/OIDC identity provider, passed as external_access_token / identity_pool_id (see BYOIDC bearer-token authentication below).

    A Global key works against every route this driver touches, so it is the more future-proof choice; a Flink Region key works against the Flink SQL routes that are the driver's focus today. Provide at least one pair. If you supply both, the Global pair is used and the Flink pair is ignored. A half-supplied pair (a key without its secret, or vice versa) is rejected.

  • Organization ID — required, unless you're using a Global API key and it can see exactly one organization, in which case you can omit organization_id and it's inferred automatically on first use of the connection. A Flink Region key has no way to discover it, so it always requires organization_id explicitly.

How to Obtain a Flink Region API Key

A Flink Region API key (also called a Flink SQL API key) is specific to your Flink region/environment (the environment + cloud provider + cloud region triplet) and provides access to the regional Flink SQL API endpoints. This is distinct from a Confluent Cloud control-plane API key. (A "Global" Confluent Cloud API key is the alternative — see above — and is created the same way but without scoping to a single Flink region.)

To create or find a Flink Region API key:

  1. Go to https://confluent.cloud/settings/api-keys
  2. Filter by resource 'Flink Region'
  3. Find an existing key, or follow the '+ Add API Key' path to create a new one
  4. When creating a new key, select either:
    • My account - for development/testing
    • Service account - for production applications (recommended) and then:
    • The Environment, Cloud Provider and Cloud Region matching the Flink database(s) / Kafka cluster(s) you intend to use this driver against.
  5. Save both the API key and API secret securely (the secret cannot be retrieved later)
  6. Use these credentials as the flink_api_key and flink_api_secret parameters in the connect() function

API keys may also be generated using the Confluent CLI or by API access, outside the scope of this document.

BYOIDC bearer-token authentication

If your Confluent Cloud organization is configured with your own registered OAuth/OIDC identity provider, you can authenticate to the Flink data plane with a bearer token you mint yourself (an "external" token in Confluent's authorization vocabulary) rather than a Confluent API key. Pass the token and the identity-pool id that scopes it:

import os
import confluent_sql

connection = confluent_sql.connect(
    external_access_token=os.environ["MY_IDP_BEARER_TOKEN"],
    identity_pool_id="pool-abc123",
    environment_id="env-...",
    organization_id="org-...",   # required under BYOIDC (see below)
    cloud_provider="aws",
    cloud_region="us-east-2",
    database="your-database-name",   # optional, and works under BYOIDC: it only sets the default
                                     # database, which needs no CMK cluster-id lookup
)

The driver stamps Authorization: Bearer <external_access_token> and Confluent-Identity-Pool-Id: <identity_pool_id> on every Flink request. The parameter name mirrors the Confluent Flink Table API plugin's client.oauth.external-access-token option. A runnable example is in examples/byoidc_bearer_token_example.py.

Things to know:

  • external_access_token and identity_pool_id are a pair and are mutually exclusive with every API-key parameter (global_*, flink_*, tableflow_*, connect_*). Supplying a bearer token alongside any API key is rejected.
  • Flink data plane only. Confluent's authorization model does not accept a raw external token on any control-plane route this driver calls. So under BYOIDC the control-plane surfaces — Tableflow (enable_tableflow / get_tableflow / disable_tableflow), Connectors, and the CMK cluster-id lookup that resolves a database name to its lkc-… id — fail closed with a clear error. Use an API-key connection for those, and pass database_kafka_cluster_id to connect() if you need the driver to know the cluster id without the CMK lookup.
  • organization_id is required under BYOIDC. Inferring it needs a control-plane call an external token can't make, so it must be supplied explicitly.
  • No token refresh. The token is used verbatim on every request. When it expires, requests begin failing (surfaced as OperationalError) and the connection is effectively dead — open a fresh connection with a fresh token.

Installation

# Using pip
pip install confluent-sql

# Using uv (recommended)
uv add confluent-sql

Quick Start

Setup the connection:

import confluent_sql

# Connect to Confluent Cloud Flink SQL
connection = confluent_sql.connect(
    organization_id="your-org-uuid",  # optional if using a Global key -- see Prerequisites
    environment_id="env-123456",
    cloud_provider="aws",
    cloud_region="us-east-2",
    flink_api_key="your-flink-api-key",  # or global_api_key=... for a "Global" Confluent Cloud key
    flink_api_secret="your-flink-api-secret",  # or global_api_secret=...
    database="your-database-name",
    compute_pool_id="lfcp-789012"  # optional; omit to use the environment's default compute pool
)

Create a cursor and run a point-in-time SNAPSHOT query:

cursor = connection.cursor()
cursor.execute("SELECT customer_id, name FROM customers")

Fetch results using fetchone(), fetchmany() and fetchall():

print(cursor.fetchone())
print(cursor.fetchmany(2))
print(cursor.fetchall())

Fetch results using the cursor as an iterator:

for row in cursor:
    print(row)

Clean up:

cursor.close()
connection.close()

Dictionary Result Rows:

...
cursor = connection.cursor(as_dict=True)
cursor.execute("SELECT customer_id, name, email FROM customers WHERE customer_id = %s", (123,))
row = cursor.fetchone()
print(row["customer_id"])  # Access by column name

Streaming Queries:

import time

...

# Execute a streaming statement, runs and produces results indefinitely until
# we stop consuming its results or the statement is stopped or deleted via API ...
cursor = connection.streaming_cursor()
cursor.execute("SELECT * FROM orders_stream WHERE total > %s", (1000,))

while cursor.may_have_results:
    rows = cursor.fetchmany(10)
    if rows:
        for row in rows:
            print(row)
    else:
        time.sleep(0.1)

Parameterized Statement and Flink to Python Value Support

This driver supports all Flink types, some with caveats. Please consult the type support documentation for more details.

DB-API Extensions

This driver extends the standard DB-API v2 interface with additional features:

  • Dictionary result rows - Access columns by name instead of position
  • Streaming cursors - Non-blocking result consumption from continuous queries
  • Changelog compression - Automatic state management for aggregations and joins
  • Statement lifecycle management - Named statements, labels, and resource management
  • Statement properties - Execution controls not expressible inline within the SQL statement
  • Type system - Full support for all Flink SQL types including streaming-specific types
  • Performance monitoring - Built-in fetch metrics and introspection

Architecture & How It Works

The confluent-sql driver communicates with Confluent Cloud Flink SQL through HTTP-based APIs. Unlike in traditional databases, statements are first-class entities on the server with their own lifecycle, allowing features like:

  • Named statements - Identify and recover queries across connections
  • Persistent execution - Statements survive connection close and can be resumed
  • Batch management - Label related statements for group operations

For an in-depth explanation of the HTTP architecture and statement lifecycle, see ARCHITECTURE.md.

Complete Documentation

For comprehensive documentation of all DB-API extensions, see DBAPI_EXTENSIONS.md.

For detailed streaming query guidance, see STREAMING.md.

For type support and examples, see TYPES.md.

Private Networking Considerations

By default, this driver uses the public Confluent Cloud API networking endpoint for the provided cloud provider and region. However, if the Flink database / Kafka cluster you intend to query requires private networking connectivity, then provide the appropriate Flink private networking base URL as the endpoint parameter to connect() or Connection.__init__(). Refer to the Confluent Cloud Flink private networking documentation for more information on composing your endpoint URL.

Symptoms of using the public endpoint when private networking is required include:

  • HTTP 429-related exceptions raised when submitting statements querying tables whose backing Kafka topics / clusters are configured for private networking only.
  • Empty or surprisingly missing results when querying INFORMATION_SCHEMA or SHOW TABLES, due to silent filtering of private-networking-only tables/topics when querying the system catalog.

Development

Setup

# Clone repository
git clone <repository-url>
cd confluent-sql

# Install uv if needed
pip install uv

# Install dependencies
uv sync

# Install in development mode
uv pip install -e .

Running Tests

Set required environment variables for integration tests. If any of the variables is not set, integration tests will be skipped.

export CONFLUENT_ORG_ID="org-123456" # Optional if using a Global key that can see exactly one org.
export CONFLUENT_ENV_ID="env-123456"
export CONFLUENT_CLOUD_PROVIDER="aws"
export CONFLUENT_CLOUD_REGION="us-east-2"
export CONFLUENT_FLINK_API_KEY="your-key" # Flink Region API key for the above cloud/region ...
export CONFLUENT_FLINK_API_SECRET="your-secret" # and associated secret.
# Alternatively, a "Global" Confluent Cloud API key (used in preference to the Flink pair if both
# are set): export CONFLUENT_GLOBAL_API_KEY / CONFLUENT_GLOBAL_API_SECRET instead. With a Global
# key, CONFLUENT_ORG_ID may be omitted -- the integration suite infers it, same as connect() does.
export CONFLUENT_TEST_DBNAME="test-db" # A database/kafka cluster name within the above cloud/region.
export CONFLUENT_COMPUTE_POOL_ID="lfcp-789012" # Optional. If set, the integration suite runs against this pool; if unset, the suite runs against the environment's default pool. The driver treats it as optional at connect() either way.

Run tests:

uv run pytest

Metadata

Release files for confluent-sql 0.5.2

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for confluent-sql 0.5.2
File Size Uploaded
confluent_sql-0.5.2.tar.gz 301.8 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for confluent-sql 0.5.2
File Interpreter ABI Platform
confluent_sql-0.5.2-py3-none-any.whl Python 3 none any Details

Total release size: 415.2 kB

Release files / confluent_sql-0.5.2.tar.gz

Download URL confluent_sql-0.5.2.tar.gz
Size 301.8 kB
Tags Source
SHA-256 checksum
How to use checksums
e3e24101e7f950db850ba7c717fd1adc73ad088daa2046e8c9d2e319290d59f6
BLAKE2b-256 checksum
How to use checksums
b576a707b28c890e0317c45c7b5b5d43653c4dfc765559bb449eec5354241f3d
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via uv/0.10.11 {"installer":{"name":"uv","version":"0.10.11","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}

Release files / confluent_sql-0.5.2-py3-none-any.whl

Download URL confluent_sql-0.5.2-py3-none-any.whl
Size 113.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
21ce9661347268c2e1fe953367702cd7710f8571c190603fca6d847e6f80d989
BLAKE2b-256 checksum
How to use checksums
b5502d17025fa758df18066a41fddbecad2f658ed4a6d71aa6dd81ba1d9100e8
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via uv/0.10.11 {"installer":{"name":"uv","version":"0.10.11","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}

Release history Release notifications | RSS feed

0.6.0

2 release files

0.5.5

2 release files

0.5.4

2 release files

0.5.3

2 release files

This release

0.5.2 This release

2 release files

0.5.1

2 release files

0.5.0

2 release files

0.4.2

2 release files

0.4.1

2 release files

0.4.0

2 release files

0.3.1

2 release files

0.3.0

2 release files

0.2.0

2 release files

0.1.3

2 release files

0.1.2

2 release files

0.1.1

2 release files

0.1.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page