Skip to content
Merged
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
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -117,8 +117,8 @@ dialect or backend you use:
| Local/S3 Iceberg and Lance | Ray Data / Daft public readers | `tributo[data,data-daft]` | Alpha; real dual-engine Conformance |
| PostgreSQL structured table | Ray Data / Daft SQL readers | `tributo[postgresql,data-daft]` | Alpha; real PostgreSQL Conformance |
| HDFS Parquet/CSV | Ray Data + PyArrow Hadoop filesystem | Ray runtime with HDFS libraries | Adapter present; cluster gate pending |
| ClickHouse | independent `daft-olap-connectors` | external package | Adapter present; package/infrastructure gates pending |
| Doris | independent `ray-doris` / `daft-olap-connectors` | external packages | Adapters present; package/infrastructure gates pending |
| ClickHouse | independent local `daft-clickhouse` wheel | `tributo[clickhouse]` plus the connector wheel | Adapter present; package/infrastructure gates pending |
| Doris | independent `ray-doris` / local `daft-doris` wheel | `tributo[mysql]` or `tributo[doris-flight]` plus the connector wheel | Adapters present; package/infrastructure gates pending |
| ORC / Hive external tables | no locked public reader path | — | Unsupported |

Provider/binding presence is not a support claim. ClickHouse, Doris, HDFS, and Hive are
Expand Down
2 changes: 1 addition & 1 deletion docs/architecture/call-chain-inventory.md
Original file line number Diff line number Diff line change
Expand Up @@ -279,7 +279,7 @@ credential-safe descriptor validation and atomic registration
WriteBinding selection; native dependency import occurs at factory/execute
```

Selected optional integrations (`ray-doris`, `daft-olap-connectors`) also have
Selected optional integrations (`ray-doris`, `daft-doris`, `daft-clickhouse`) also have
thin built-in descriptors and explicit install diagnostics. Their adapters are
not support claims until their external packages and infrastructure gates pass.

Expand Down
4 changes: 2 additions & 2 deletions docs/data/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -56,8 +56,8 @@ selects `binding_id` explicitly.
| Local/S3 Lance | Native reader | Native reader | Verified |
| PostgreSQL structured table | Native SQL reader | Native SQL reader | Verified |
| HDFS Parquet/CSV | Native reader with PyArrow HDFS | No locked public reader | Adapted; cluster gate pending |
| ClickHouse | No selected Binding | `daft-olap-connectors` | Adapter only; external package and database gates pending |
| Doris | `ray-doris` | `daft-olap-connectors` | Adapter only; external packages and database gates pending |
| ClickHouse | No selected Binding | External `daft-clickhouse` wheel | Adapter only; external package and database gates pending |
| Doris | `ray-doris` | External `daft-doris` wheel | Adapter only; external package and database gates pending |
| ORC or Hive external table | No locked public reader | No locked public reader | Unsupported, fail-closed |

“Verified” means the current combination has semantic Conformance and real
Expand Down
11 changes: 11 additions & 0 deletions docs/reference/api/data.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,17 @@ a public annotation or moving a public object.
Stable, Beta, and Alpha objects appear because Ray-style API policy requires
documentation for every public stability tier.

## `tributo.data.bindings._daft_sql`

```{autoclass} tributo.data.bindings._daft_sql.DaftClickHouseBinding
:no-members:
```

```{autoclass} tributo.data.bindings._daft_sql.DaftDorisBinding
:no-members:
```


## `tributo.data.contracts.handles`

```{autoclass} tributo.data.contracts.handles.DaftDataFrameHandle
Expand Down
4 changes: 2 additions & 2 deletions docs/reference/support-matrix.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,8 @@ compatible profile, while the generated `Validated profiles` column remains
| PostgreSQL structured table reads | Verified | Ray uses a single public SQL read and fails closed on parallel shard requirements; Daft may use native partition hints |
| ClickHouse/Doris raw SQL | Unsupported | Legacy shapes return a credential-free migration error; use structured table input or execute SQL outside Tributo ingestion |
| HDFS Parquet/CSV reads | Adapter only | Ray binding exists; real HDFS/JVM/worker gate is pending |
| ClickHouse reads | Adapter only | Requires unpublished `daft-olap-connectors` and real-database Conformance; provider partition discovery is distinct from engine auto-routing |
| Doris reads | Adapter only | Requires unpublished `ray-doris` or `daft-olap-connectors` and real-database Conformance; tablet planning remains provider/binding-owned |
| ClickHouse reads | Adapter only | Requires an externally installed `daft-clickhouse` wheel and real-database Conformance; provider partition discovery is distinct from engine auto-routing |
| Doris reads | Adapter only | Requires `ray-doris` or an externally installed `daft-doris` wheel and real-database Conformance; tablet planning remains provider/binding-owned |
| ORC and Hive external-table reads | Not implemented | Locked Ray/Daft versions expose no validated public reader |
| Third-party ingestion Provider/Binding SPI | Implemented | Installed packages use `tributo.ingestion_providers` plus `tributo.ingestion_bindings`; bad plugins are isolated, duplicate routes never replace built-ins, and Binding selection can constrain filesystem, catalog, and storage format |
| Lance output | Implemented as a generic ResultSink path | User Predictor owns vector semantics; the sink does not pool, normalize, or automatically invoke the separate vector-index workflow |
Expand Down
15 changes: 10 additions & 5 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@ data = [
"pylance==9.0.0",
]
data-daft = [
"daft[lance,ray]>=0.7.0,<0.8.0",
"daft[lance,ray]>=0.7.23,<0.7.24",
]
vector-index = [
"tributo[data,s3]",
Expand All @@ -71,10 +71,15 @@ postgresql = [
"sqlglot<30.9.0",
]
clickhouse = [
"clickhouse-connect>=1.3.0",
"clickhouse-connect[arrow,async]>=1.5,<1.6",
]
mysql = [
"pymysql>=1.1.0",
"PyMySQL>=1.2,<1.3",
]
doris-flight = [
"PyMySQL>=1.2,<1.3",
"adbc-driver-manager>=1.6,<2",
"adbc-driver-flightsql>=1.6,<2",
]
s3 = [
"boto3>=1.42.91",
Expand Down Expand Up @@ -336,9 +341,9 @@ dev = [
"pytest-asyncio>=1.3.0",
"pytest-metadata>=3.1.1",
"redis>=8.0.0",
"clickhouse-connect>=1.3.0",
"mypy>=2.3.0",
"daft>=0.7.0,<0.8.0",
"clickhouse-connect[arrow,async]>=1.5,<1.6",
"daft>=0.7.23,<0.7.24",
"ml-dtypes>=0.5.0",
"skl2onnx>=1.17.0",
]
103 changes: 68 additions & 35 deletions src/tributo/data/bindings/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,11 +22,19 @@
_DEFAULT_BINDINGS_LOCK = threading.Lock()
_RAY_VERSION_SPEC = "==2.55.1"
_DAFT_VERSION_SPEC = ">=0.7.0,<0.8.0"
_DAFT_SQL_VERSION_SPEC = ">=0.7.23,<0.7.24"
_RAY_INSTALL_HINT = "pip install 'ray[default,serve,tune]==2.55.1'"
_DAFT_INSTALL_HINT = "pip install 'tributo[data-daft]'"
_DAFT_LANCE_INSTALL_HINT = "pip install 'tributo[data,data-daft]'"
_DATA_INSTALL_HINT = "pip install 'tributo[data]'"
_DAFT_OLAP_INSTALL_HINT = "pip install 'daft-olap-connectors[clickhouse,doris]'"
_DAFT_CLICKHOUSE_INSTALL_HINT = (
"Install a local daft-clickhouse[clickhouse] wheel, then run "
"pip install 'tributo[clickhouse]'"
)
_DAFT_DORIS_INSTALL_HINT = (
"Install a local daft-doris[doris] wheel (or daft-doris[doris-flight] for Flight), "
"then run pip install 'tributo[mysql]' (or 'tributo[doris-flight]' for Flight)"
)
_RAY_DORIS_INSTALL_HINT = "pip install 'ray-doris[mysql,flight]'"
_POSTGRESQL_INSTALL_HINT = "pip install 'tributo[postgresql]'"

Expand Down Expand Up @@ -262,33 +270,49 @@ def _daft_lance_descriptor() -> BindingDescriptor:
)


def _daft_olap_descriptor(connector_id: str) -> BindingDescriptor:
from tributo.data.bindings.daft_olap import (
DaftClickHouseBinding,
DaftDorisBinding,
)
def _daft_clickhouse_descriptor() -> BindingDescriptor:
from tributo.data.bindings.daft_clickhouse import DaftClickHouseBinding

factory = (
DaftClickHouseBinding if connector_id == "clickhouse" else DaftDorisBinding
)
return BindingDescriptor(
key=BindingKey(
"tributo.daft",
ScanKind.SQL,
connector_id,
f"daft_olap.daft.{connector_id}",
"clickhouse",
"daft_clickhouse.daft.clickhouse",
),
factory=factory,
factory=DaftClickHouseBinding,
capabilities=frozenset({SourceCapability.PROJECTION}),
distribution_name="daft-olap-connectors",
distribution_version=(
_distribution_version("daft-olap-connectors") or "0.1.0a1"
distribution_name="daft-clickhouse",
distribution_version=_distribution_version("daft-clickhouse") or "0.1.0a1",
engine_version_spec=_DAFT_SQL_VERSION_SPEC,
dependency_distributions=("clickhouse-connect",),
supported_read_hints=frozenset(
{ReadHint.TARGET_PARALLELISM, ReadHint.BATCH_SIZE}
),
engine_version_spec=_DAFT_VERSION_SPEC,
install_hint=_DAFT_CLICKHOUSE_INSTALL_HINT,
)


def _daft_doris_descriptor() -> BindingDescriptor:
from tributo.data.bindings.daft_doris import DaftDorisBinding

return BindingDescriptor(
key=BindingKey(
"tributo.daft",
ScanKind.SQL,
"doris",
"daft_doris.daft.doris",
),
factory=DaftDorisBinding,
capabilities=frozenset({SourceCapability.PROJECTION}),
distribution_name="daft-doris",
distribution_version=_distribution_version("daft-doris") or "0.1.0a1",
engine_version_spec=_DAFT_SQL_VERSION_SPEC,
dependency_distributions=("PyMySQL",),
supported_read_hints=frozenset(
{ReadHint.TARGET_PARALLELISM, ReadHint.BATCH_SIZE}
),
install_hint=_DAFT_OLAP_INSTALL_HINT,
install_hint=_DAFT_DORIS_INSTALL_HINT,
)


Expand Down Expand Up @@ -570,24 +594,33 @@ def default_engine_bindings() -> EngineBindings:
None,
("pylance", "daft-lance"),
),
*(
(
lambda connector_id=connector_id: _daft_olap_descriptor(
connector_id
),
BindingKey(
"tributo.daft",
ScanKind.SQL,
connector_id,
f"daft_olap.daft.{connector_id}",
),
"daft",
_DAFT_VERSION_SPEC,
_DAFT_OLAP_INSTALL_HINT,
None,
("daft-olap-connectors",),
)
for connector_id in ("clickhouse", "doris")
(
_daft_clickhouse_descriptor,
BindingKey(
"tributo.daft",
ScanKind.SQL,
"clickhouse",
"daft_clickhouse.daft.clickhouse",
),
"daft",
_DAFT_SQL_VERSION_SPEC,
_DAFT_CLICKHOUSE_INSTALL_HINT,
None,
("daft-clickhouse", "clickhouse-connect"),
),
(
_daft_doris_descriptor,
BindingKey(
"tributo.daft",
ScanKind.SQL,
"doris",
"daft_doris.daft.doris",
),
"daft",
_DAFT_SQL_VERSION_SPEC,
_DAFT_DORIS_INSTALL_HINT,
None,
("daft-doris", "PyMySQL"),
),
(
_ray_doris_descriptor,
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
"""Daft OLAP Bindings delegating to the independent connector package."""
"""Shared Daft SQL bindings for independent database connector packages."""

from __future__ import annotations

Expand Down Expand Up @@ -30,24 +30,37 @@
apply_pipeline_to_daft_df,
)
from tributo.exceptions import JobConfigurationError
from tributo.util.annotations import PublicAPI, Stability


@dataclass(frozen=True)
class _DaftOlapNativePlan:
class _DaftSqlNativePlan:
dataframe: Any
input_schema: pa.Schema
transforms: CompiledPipeline
reader_api: str
transport_id: str


class _DaftOlapBinding:
class _DaftSqlBinding:
connector_id: ClassVar[str]
package_name: ClassVar[str]
reader_api: ClassVar[str]

def compile(self, request: BindingCompileRequest) -> BindingCompilation:
with binding_stage("validate_capabilities"):
plan = require_sql_table(request.plan, self.connector_id)
if plan.sharding.mode is SqlShardMode.PARALLEL:
raise BindingStageError.framework_diagnostic(
"validate_capabilities",
error_type=JobConfigurationError,
diagnostic_code="parallel_sql_read_unsupported",
diagnostic=(
"The Daft SQL Binding supports partitioning.mode "
"'single' or 'auto'; explicit 'parallel' with shard "
"columns is unsupported"
),
)
if (
plan.sharding.mode is SqlShardMode.SINGLE
and request.read_options.target_parallelism is not None
Expand All @@ -58,10 +71,19 @@ def compile(self, request: BindingCompileRequest) -> BindingCompilation:
diagnostic_code="single_sql_read_rejects_parallelism_hint",
diagnostic=(
"A single SQL read cannot honor target_parallelism; "
"set partitioning.mode to 'auto' or 'parallel', or "
"remove target_parallelism"
"set partitioning.mode to 'auto' or remove "
"target_parallelism"
),
)
if self.connector_id == "doris":
protocol = str(request.runtime_options.get("protocol") or "mysql")
if protocol not in {"mysql", "flight"}:
raise BindingStageError.framework_diagnostic(
"validate_capabilities",
error_type=JobConfigurationError,
diagnostic_code="unsupported_doris_transport",
diagnostic=("Doris transport must be 'mysql' or 'flight'"),
)
with binding_stage("classify_transforms"):
decisions = residual_decisions(request.transforms)
with binding_stage("build_native_plan"):
Expand All @@ -71,9 +93,7 @@ def compile(self, request: BindingCompileRequest) -> BindingCompilation:

def _build(
self, request: BindingCompileRequest, plan: SqlScan
) -> _DaftOlapNativePlan:
from daft_olap import read_clickhouse, read_doris

) -> _DaftSqlNativePlan:
target = resolve_sql_target(plan, request.runtime_options)
options: dict[str, Any] = {
"host": target.host,
Expand All @@ -93,11 +113,16 @@ def _build(
)
if target_tasks is not None:
options["target_tasks"] = target_tasks

if self.connector_id == "clickhouse":
from daft_clickhouse import read_clickhouse

options["port"] = target.port
dataframe = read_clickhouse(**options)
transport_id = "clickhouse_native"
else:
from daft_doris import read_doris

protocol = str(request.runtime_options.get("protocol") or "mysql")
options["transport"] = protocol
options["mysql_port"] = target.port
Expand All @@ -109,21 +134,22 @@ def _build(
options["flight_port"] = int(flight_port)
dataframe = read_doris(**options)
transport_id = protocol

schema = canonical_engine_schema(dataframe.schema())
transforms = ConcreteTransformCompiler().compile(
request.transforms, TransformBackend.DAFT, schema
)
return _DaftOlapNativePlan(
return _DaftSqlNativePlan(
dataframe,
schema,
transforms,
self.reader_api,
transport_id,
)

@staticmethod
def _wrap(
native_plan: _DaftOlapNativePlan,
self,
native_plan: _DaftSqlNativePlan,
decisions: tuple[TransformDecision, ...],
) -> BindingCompilation:
transformed = apply_pipeline_to_daft_df(
Expand All @@ -144,17 +170,23 @@ def _wrap(
schema_fingerprint=schema_fingerprint(output_schema),
metadata_fetched=True,
physical_splits=PhysicalSplitSummary(
detail="database splits and batches are delegated to daft-olap-connectors"
detail=(
f"database splits and batches are delegated to {self.package_name}"
)
),
diagnostics=("database metadata I/O was used for schema inference",),
)


class DaftClickHouseBinding(_DaftOlapBinding):
@PublicAPI(stability=Stability.ALPHA)
class DaftClickHouseBinding(_DaftSqlBinding):
connector_id = "clickhouse"
reader_api = "daft_olap.read_clickhouse"
package_name = "daft-clickhouse"
reader_api = "daft_clickhouse.read_clickhouse"


class DaftDorisBinding(_DaftOlapBinding):
@PublicAPI(stability=Stability.ALPHA)
class DaftDorisBinding(_DaftSqlBinding):
connector_id = "doris"
reader_api = "daft_olap.read_doris"
package_name = "daft-doris"
reader_api = "daft_doris.read_doris"
5 changes: 5 additions & 0 deletions src/tributo/data/bindings/daft_clickhouse.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
"""Daft ClickHouse Binding using the public daft-clickhouse facade."""

from tributo.data.bindings._daft_sql import DaftClickHouseBinding

__all__ = ["DaftClickHouseBinding"]
Loading
Loading