| 文件 | 最后提交记录 | 最后更新时间 |
|---|---|---|
feat: Add Catalog and Table implementations for PostgreSQL (#5487) ## Changes Made Allows us to do catalog = daft.catalog.Catalog.from_postgres("postgresql://user:pass@uri:5432/postgres") table = catalog.get_table("myschema.mytable") df = table.read() ... output_table = catalog.create_table_if_not_exists("myschema.output") output_table.write(df) Depends on https://github.com/Eventual-Inc/Daft/pull/5471 for distributed writes and the df.write_sql() interface. ## Key decisions, limitations, and notes - Embeddings support is currently provided via the pgvector extension and the pgvector python client library. - The following types are currently not supported: - Duration - Interval - The following types currently fallback to JSONB: - Multidimensional lists - Structs - Maps - Images - Tensors - No connection pooling - Writes happen on a single node (https://github.com/Eventual-Inc/Daft/pull/5471 implements distributed writes) - Uses psycopg and COPY for writes. ADBC would potentially be a better choice here but we ignore it for now due to its limited type support. - Instead of assuming that public is the default schema, we delegate table resolution to Postgres' [Schema Search Path](https://www.postgresql.org/docs/current/ddl-schemas.html#DDL-SCHEMAS-PATH) --------- Co-authored-by: greptile-apps[bot] <165735046+greptile-apps[bot]@users.noreply.github.com> | 10 个月前 | |
[FEAT] read_sql (#1943) Closes #1560 Adds a new read method: read_sql(sql: str, url: str), which executes a given sql query on a given database url, and creates a Dataframe from the results. Drive bys: - Added a from_arrow cast from pyarrow date64 to daft timestamp (both are millisecond precision) to allow connectorx reads with timestamp columns. Features: - Uses [connector-x](https://github.com/sfu-db/connector-x/tree/main), which reads directly into pyarrow via rust when supported, else fallback to [SQL alchemy ](https://docs.sqlalchemy.org/en/20/orm/quickstart.html) - Partitioned reads via limit and offset when supported - Integration tests for Postgres, MySQL, Trino - Pushdowns into base SQL query | 2 年前 | |
chore: Upgrade Ruff ruleset to 3.9 and add from __future__ import annotations (#4393) | 1 年前 | |
fix: Map Literal <-> Python Dict conversion (#6084) ## Changes Made Return map literals as Python dicts instead of lists of tuples. Update tests to expect dicts for map outputs and value_counts results. ## Related Issues #6081 --------- Co-authored-by: desmondcheongzx <desmondcheongzx@gmail.com> Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com> | 7 个月前 | |
chore: Remove the old Ray Runner (#5375) ## Changes Made 🎉🎂🥳 Can finally delete it, Flotilla supports all necessary features. Additional Related Features Removed: * DataFrame.num_partitions: This is computed using the old Ray runner's planner. We could use the new Ray runner, but it seems kind of unnecessary * Context Settings: --------- Co-authored-by: Colin Ho <colin.ho99@gmail.com> | 10 个月前 | |
[CHORE] Add tpch test for read sql (#2026) This PR adds tpch tests via read_sql. Our tpch integration test setup already creates a sqlite table, which makes it easy to run read_sql against for complex query validation. As a drive by in addition to adding tpch tests, this PR gates the Date and Timestamp literal to sql translations. This is because the semantics for date and timestamp literals were not as general as expected. | 2 年前 | |
feat(io): implement write_sql with SQLDataSink and explicit dtype support (#5979) ## Summary This PR introduces DataFrame.write_sql() , enabling users to write Daft DataFrames to SQL databases (e.g., PostgreSQL, SQLite) via SQLAlchemy. It implements a robust, distributed SQLDataSink that: 1. Handles Distributed Writes : Uses the DataSink pattern to manage driver-side table initialization and worker-side parallel writes. 2. Supports Explicit Types : Exposes a dtype parameter to allow users to override default type inference with specific SQLAlchemy types (addressing type verification concerns). 3. Ensures Connection Safety : Manages connection lifecycles properly across distributed workers to avoid socket serialization issues. ## Key Changes ### 1. New Public API: DataFrame.write_sql - Location : daft/dataframe/dataframe.py - Signature : def write_sql( self, table_name: str, conn: str | Callable[[], "Connection"], write_mode: Literal["append", "overwrite", "fail"] = "append", chunk_size: int | None = None, dtype: dict[str, Any] | None = None # <--- NEW: Explicit type control ) -> DataFrame: ... - Behavior : Delegates to SQLDataSink and returns a DataFrame with write metrics ( total_written_rows , total_written_bytes ). ### 2. Internal Implementation: SQLDataSink - Location : daft/io/_sql.py - Architecture : - start() (Driver) : Handles write_mode logic. - fail : Checks existence, raises error if table exists, creates table schema if not. - overwrite : Replaces table with new schema. - append : Creates table if not exists, ensuring schema readiness for workers. - write() (Workers) : - Creates independent SQLAlchemy engines/connections per task (no pickling of connections). - Converts MicroPartitions to Pandas. - Uses pd.to_sql with the user-provided dtype to write data efficiently. - Ensures proper resource cleanup ( engine.dispose() ). ### 3. Tests - Location : tests/integration/sql/test_write_sql.py - Coverage : - Sources : Verified with PyDict, CSV, and JSON sources. - Modes : Comprehensive tests for append , overwrite , and fail modes. - Type Verification : Added specific tests ( test_write_sql_dtype_basic_types ) that use sqlalchemy.inspect to verify that columns are created with the correct SQL types when dtype is provided. - Connection Factory : Verified support for passing a connection factory function (crucial for pickling compatibility). ## Addressing Previous Concerns (Type Verification) This implementation addresses concerns about type safety (raised in previous discussions) by: 1. Leveraging Pandas' mature type inference for standard types. 2. Providing the dtype "escape hatch" for complex scenarios, giving users full control over the target schema definition. 3. Including integration tests that explicitly verify schema creation correctness. ## Checklist - I have added comprehensive unit/integration tests. - I have updated the documentation (docstrings included). - I have verified that connections are properly closed and disposed of. Thank you very much for the idea provided by https://github.com/Eventual-Inc/Daft/pull/5471 ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> | 7 个月前 |
| 文件 | 最后提交记录 | 最后更新时间 |
|---|---|---|
| 10 个月前 | ||
| 2 年前 | ||
| 1 年前 | ||
| 7 个月前 | ||
| 10 个月前 | ||
| 2 年前 | ||
| 7 个月前 |