Skip to content

fsspeckit.datasets.duckdb API Reference

duckdb

DuckDB dataset integration for fsspeckit.

This package contains focused submodules for DuckDB functionality: - dataset: Dataset I/O and maintenance operations - connection: Connection management and filesystem registration - helpers: Utility functions for DuckDB operations

All public APIs are re-exported here for convenient access.

Classes

fsspeckit.datasets.duckdb.DuckDBConnection

DuckDBConnection(
    filesystem: AbstractFileSystem | None = None,
)

Manages DuckDB connection lifecycle and filesystem registration.

This class is responsible for: - Creating and managing DuckDB connections - Registering fsspec filesystems with DuckDB - Connection cleanup

Parameters:

Name Type Description Default
filesystem AbstractFileSystem | None

fsspec filesystem instance to use

None

Initialize DuckDB connection manager.

Parameters:

Name Type Description Default
filesystem AbstractFileSystem | None

Filesystem to use. Defaults to local filesystem.

None
Source code in src/fsspeckit/datasets/duckdb/connection.py
def __init__(self, filesystem: AbstractFileSystem | None = None) -> None:
    """Initialize DuckDB connection manager.

    Args:
        filesystem: Filesystem to use. Defaults to local filesystem.
    """
    self._connection: duckdb.DuckDBPyConnection | None = None
    self._filesystem = filesystem or fsspec_filesystem("file")
Attributes
fsspeckit.datasets.duckdb.DuckDBConnection.connection property
connection: Any

Get active DuckDB connection, creating it if necessary.

Returns:

Type Description
Any

Active DuckDB connection

fsspeckit.datasets.duckdb.DuckDBConnection.filesystem property
filesystem: AbstractFileSystem

Get the filesystem instance.

Returns:

Type Description
AbstractFileSystem

Filesystem instance

Functions
fsspeckit.datasets.duckdb.DuckDBConnection.__del__
__del__() -> None

Destructor to ensure connection is closed.

Source code in src/fsspeckit/datasets/duckdb/connection.py
def __del__(self) -> None:
    """Destructor to ensure connection is closed."""
    self.close()
fsspeckit.datasets.duckdb.DuckDBConnection.__enter__
__enter__() -> DuckDBConnection

Enter context manager.

Returns:

Type Description
DuckDBConnection

self

Source code in src/fsspeckit/datasets/duckdb/connection.py
def __enter__(self) -> DuckDBConnection:
    """Enter context manager.

    Returns:
        self
    """
    return self
fsspeckit.datasets.duckdb.DuckDBConnection.__exit__
__exit__(
    exc_type: type[BaseException] | None,
    exc_value: BaseException | None,
    traceback: Any,
) -> None

Exit context manager and close connection.

Source code in src/fsspeckit/datasets/duckdb/connection.py
def __exit__(
    self,
    exc_type: type[BaseException] | None,
    exc_value: BaseException | None,
    traceback: Any,
) -> None:
    """Exit context manager and close connection."""
    self.close()
fsspeckit.datasets.duckdb.DuckDBConnection.close
close() -> None

Close the connection and clean up resources.

Source code in src/fsspeckit/datasets/duckdb/connection.py
def close(self) -> None:
    """Close the connection and clean up resources."""
    if self._connection is not None:
        try:
            self._connection.close()
        except (ConnectionException, OperationalError) as e:
            logger.warning("Error closing DuckDB connection: %s", e)
        finally:
            self._connection = None
fsspeckit.datasets.duckdb.DuckDBConnection.execute_sql
execute_sql(
    query: str, parameters: list[Any] | None = None
) -> Any

Execute a SQL query.

Parameters:

Name Type Description Default
query str

SQL query to execute

required
parameters list[Any] | None

Optional query parameters

None

Returns:

Type Description
Any

Query result

Source code in src/fsspeckit/datasets/duckdb/connection.py
def execute_sql(
    self,
    query: str,
    parameters: list[Any] | None = None,
) -> Any:
    """Execute a SQL query.

    Args:
        query: SQL query to execute
        parameters: Optional query parameters

    Returns:
        Query result
    """
    conn = self.connection

    if parameters:
        return conn.execute(query, parameters)
    else:
        return conn.execute(query)

fsspeckit.datasets.duckdb.DuckDBDatasetIO

DuckDBDatasetIO(connection: DuckDBConnection)

Bases: BaseDatasetHandler

DuckDB-based dataset I/O operations.

This class provides methods for reading and writing parquet files and datasets using DuckDB's high-performance parquet engine.

Inherits the BaseDatasetHandler contract to provide a consistent interface across different backend implementations.

Parameters:

Name Type Description Default
connection DuckDBConnection

DuckDB connection manager

required

Initialize DuckDB dataset I/O.

Parameters:

Name Type Description Default
connection DuckDBConnection

DuckDB connection manager

required
Source code in src/fsspeckit/datasets/duckdb/dataset.py
def __init__(self, connection: DuckDBConnection) -> None:
    """Initialize DuckDB dataset I/O.

    Args:
        connection: DuckDB connection manager
    """
    self._connection = connection
Attributes
fsspeckit.datasets.duckdb.DuckDBDatasetIO.filesystem property
filesystem: 'AbstractFileSystem'

Return the filesystem instance used by this handler.

Functions
fsspeckit.datasets.duckdb.DuckDBDatasetIO.merge
merge(
    data: Table | list[Table],
    path: str,
    strategy: Literal["insert", "update", "upsert"],
    key_columns: list[str] | str,
    *,
    partition_columns: list[str] | str | None = None,
    schema: Schema | None = None,
    compression: str | None = "snappy",
    max_rows_per_file: int | None = 5000000,
    row_group_size: int | None = 500000,
    merge_chunk_size_rows: int = 100000,
    enable_streaming_merge: bool = True,
    merge_max_memory_mb: int = 1024,
    merge_max_process_memory_mb: int | None = None,
    merge_min_system_available_mb: int = 512,
    merge_progress_callback: Callable[[int, int], None]
    | None = None,
    use_merge: bool | None = None,
) -> "MergeResult"

Merge data into an existing parquet dataset incrementally (DuckDB backend).

Semantics: - insert: append only new keys as new file(s); never rewrites existing files. - update: rewrite only files that actually contain keys being updated; never inserts. - upsert: rewrite only affected files and append inserted keys as new file(s).

Parameters:

Name Type Description Default
use_merge bool | None

Ignored (reserved for backward compatibility).

None
merge_chunk_size_rows int

Streaming merge chunk size (ignored by DuckDB).

100000
enable_streaming_merge bool

Streaming merge toggle (ignored by DuckDB).

True
merge_max_memory_mb int

Max PyArrow memory in MB (ignored by DuckDB).

1024
merge_max_process_memory_mb int | None

Max process RSS in MB (ignored by DuckDB).

None
merge_min_system_available_mb int

Min system available memory in MB (ignored by DuckDB).

512
merge_progress_callback Callable[[int, int], None] | None

Progress callback (ignored by DuckDB).

None
Source code in src/fsspeckit/datasets/duckdb/dataset.py
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
def merge(
    self,
    data: pa.Table | list[pa.Table],
    path: str,
    strategy: Literal["insert", "update", "upsert"],
    key_columns: list[str] | str,
    *,
    partition_columns: list[str] | str | None = None,
    schema: pa.Schema | None = None,
    compression: str | None = "snappy",
    max_rows_per_file: int | None = 5_000_000,
    row_group_size: int | None = 500_000,
    merge_chunk_size_rows: int = 100_000,
    enable_streaming_merge: bool = True,
    merge_max_memory_mb: int = 1024,
    merge_max_process_memory_mb: int | None = None,
    merge_min_system_available_mb: int = 512,
    merge_progress_callback: Callable[[int, int], None] | None = None,
    use_merge: bool | None = None,
) -> "MergeResult":
    """Merge data into an existing parquet dataset incrementally (DuckDB backend).

    Semantics:
    - `insert`: append only new keys as new file(s); never rewrites existing files.
    - `update`: rewrite only files that actually contain keys being updated; never inserts.
    - `upsert`: rewrite only affected files and append inserted keys as new file(s).

    Args:
        use_merge: Ignored (reserved for backward compatibility).
        merge_chunk_size_rows: Streaming merge chunk size (ignored by DuckDB).
        enable_streaming_merge: Streaming merge toggle (ignored by DuckDB).
        merge_max_memory_mb: Max PyArrow memory in MB (ignored by DuckDB).
        merge_max_process_memory_mb: Max process RSS in MB (ignored by DuckDB).
        merge_min_system_available_mb: Min system available memory in MB (ignored by DuckDB).
        merge_progress_callback: Progress callback (ignored by DuckDB).
    """
    import pyarrow.compute as pc
    import pyarrow.parquet as pq

    from fsspeckit.core.incremental import (
        IncrementalFileManager,
        MergeFileMetadata,
        MergeResult,
        confirm_affected_files,
        extract_source_partition_values,
        list_dataset_files,
        parse_hive_partition_path,
        plan_incremental_rewrite,
        validate_no_null_keys,
    )

    validate_path(path)
    validate_compression_codec(compression)
    row_group_size = self._validate_write_parameters(
        max_rows_per_file,
        row_group_size,
    )

    if use_merge is not None:
        logger.debug("duckdb_merge_use_merge_ignored", use_merge=use_merge)

    # Combine source input to a single table.
    source_table = self._combine_tables(data)

    if schema is not None:
        from fsspeckit.datasets.schema import cast_schema

        source_table = cast_schema(source_table, schema)

    key_cols = self._validate_key_columns(
        key_columns,
        source_table.column_names,
        context="source",
    )
    partition_cols = self._validate_partition_columns(
        partition_columns,
        source_table.column_names,
    )

    validate_no_null_keys(source_table, key_cols)

    fs = self._connection.filesystem

    # List existing parquet files in the dataset. Target discovery stays
    # backend-local; core planning receives backend-neutral metadata only.
    target_files = list_dataset_files(path, filesystem=fs)
    target_metadata = MergeTargetMetadata(
        exists=bool(target_files),
        files=target_files,
        row_count=sum(
            pq.read_metadata(f, filesystem=fs).num_rows for f in target_files
        ),
    )

    plan = plan_merge_operation(
        source_table=source_table,
        strategy=strategy,
        key_columns=key_cols,
        target_metadata=target_metadata,
        partition_columns=partition_cols,
    )

    source_table = plan.source_table
    key_cols = plan.key_columns
    partition_cols = plan.partition_columns
    source_keys = plan.source_keys
    source_key_set = plan.source_key_set
    target_files = plan.target_files
    target_exists = plan.target_exists
    target_count_before = plan.target_count_before

    early_result = resolve_merge_plan_early_exit(plan)
    if early_result is not None:
        return early_result

    if not target_exists:
        # INSERT/UPSERT into a non-existent dataset: write all rows as inserts.
        fs.mkdirs(path, exist_ok=True)
        write_res = self.write_dataset(
            source_table,
            path,
            mode="append",
            partition_by=partition_cols or None,
            compression=compression,
            max_rows_per_file=max_rows_per_file,
            row_group_size=row_group_size,
        )

        inserted_files = [m.path for m in write_res.files]
        files_meta = [
            MergeFileMetadata(
                path=m.path,
                row_count=m.row_count,
                operation="inserted",
                size_bytes=m.size_bytes,
            )
            for m in write_res.files
        ]

        return MergeResult(
            strategy=strategy,
            source_count=source_table.num_rows,
            target_count_before=0,
            target_count_after=write_res.total_rows,
            inserted=write_res.total_rows,
            updated=0,
            deleted=0,
            files=files_meta,
            rewritten_files=[],
            inserted_files=inserted_files,
            preserved_files=[],
        )

    # Existing dataset: plan incremental rewrite candidates using metadata.
    source_partition_values: set[tuple[object, ...]] | None = None
    if partition_cols:
        source_partition_values = extract_source_partition_values(
            source_table, partition_cols
        )

    rewrite_plan = plan_incremental_rewrite(
        dataset_path=path,
        source_keys=source_keys,
        key_columns=key_cols,
        filesystem=fs,
        partition_columns=partition_cols or None,
        source_partition_values=source_partition_values,
    )

    # Confirm actual affected files by scanning key columns.
    affected_files = confirm_affected_files(
        candidate_files=rewrite_plan.affected_files,
        key_columns=key_cols,
        source_keys=source_keys,
        filesystem=fs,
    )

    # Compute per-file matched keys for accurate updates and insert determination.
    matched_keys: set[object] = set()
    matched_keys_by_file: dict[str, set[object]] = {}
    for file_path in affected_files:
        try:
            key_table = pq.read_table(
                file_path,
                columns=key_cols,
                filesystem=fs,
                partitioning=None,
            )
            if len(key_cols) == 1:
                file_keys = set(key_table.column(key_cols[0]).to_pylist())
            else:
                file_keys = set(
                    zip(*[key_table.column(c).to_pylist() for c in key_cols])
                )
            file_matched = source_key_set & file_keys
            if file_matched:
                matched_keys_by_file[file_path] = set(file_matched)
                matched_keys |= set(file_matched)
        except (OSError, IOError, Exception):
            # Conservative: assume all source keys might be present.
            matched_keys_by_file[file_path] = set(source_key_set)
            matched_keys |= set(source_key_set)

    inserted_key_set = source_key_set - matched_keys

    if strategy == "insert":
        preserved_files = list(target_files)

        if not inserted_key_set:
            return MergeResult(
                strategy="insert",
                source_count=source_table.num_rows,
                target_count_before=target_count_before,
                target_count_after=target_count_before,
                inserted=0,
                updated=0,
                deleted=0,
                files=[
                    MergeFileMetadata(path=f, row_count=0, operation="preserved")
                    for f in preserved_files
                ],
                rewritten_files=[],
                inserted_files=[],
                preserved_files=preserved_files,
            )

        insert_table = self._select_rows_by_keys(
            source_table,
            key_cols,
            set(inserted_key_set),
        )
        write_res = self.write_dataset(
            insert_table,
            path,
            mode="append",
            partition_by=partition_cols or None,
            compression=compression,
            max_rows_per_file=max_rows_per_file,
            row_group_size=row_group_size,
        )

        inserted_files = [m.path for m in write_res.files]
        inserted_meta = [
            MergeFileMetadata(
                path=m.path,
                row_count=m.row_count,
                operation="inserted",
                size_bytes=m.size_bytes,
            )
            for m in write_res.files
        ]

        files_meta = [
            MergeFileMetadata(path=f, row_count=0, operation="preserved")
            for f in preserved_files
        ] + inserted_meta

        return MergeResult(
            strategy="insert",
            source_count=source_table.num_rows,
            target_count_before=target_count_before,
            target_count_after=target_count_before + insert_table.num_rows,
            inserted=insert_table.num_rows,
            updated=0,
            deleted=0,
            files=files_meta,
            rewritten_files=[],
            inserted_files=inserted_files,
            preserved_files=preserved_files,
        )

    # UPDATE / UPSERT: rewrite only actually affected files.
    file_manager = IncrementalFileManager()
    staging_dir = file_manager.create_staging_directory(path, filesystem=fs)

    rewritten_files: list[str] = []
    rewritten_meta: list[MergeFileMetadata] = []

    preserved_files = [f for f in target_files if f not in affected_files]

    # Prepare a match marker for join-driven full-row replacement.
    match_col_name = "__fsspeckit_match"
    if match_col_name in source_table.column_names:
        raise ValueError(f"Source contains reserved column: {match_col_name}")

    import pyarrow as pa_mod

    source_with_match = source_table.append_column(
        match_col_name, pa_mod.array([True] * source_table.num_rows)
    )

    try:
        for file_path in affected_files:
            file_matched = matched_keys_by_file.get(file_path, set())
            if not file_matched:
                preserved_files.append(file_path)
                continue

            target_table = pq.read_table(
                file_path,
                filesystem=fs,
                partitioning=None,
            )
            output_columns = target_table.column_names

            if partition_cols:
                partition_values = parse_hive_partition_path(
                    file_path,
                    partition_columns=partition_cols,
                )
                for col, value in partition_values.items():
                    if col in target_table.column_names:
                        continue
                    target_table = target_table.append_column(
                        col,
                        pa_mod.array([value] * target_table.num_rows),
                    )
            source_for_file = self._select_rows_by_keys(
                source_with_match,
                key_cols,
                set(file_matched),
            )

            joined = target_table.join(
                source_for_file,
                keys=key_cols,
                join_type="left outer",
                right_suffix="__src",
                coalesce_keys=True,
            )

            match_mask = pc.is_valid(joined.column(match_col_name))

            if partition_cols:
                for col in partition_cols:
                    if col in key_cols:
                        continue
                    src_name = f"{col}__src"
                    if src_name not in joined.column_names:
                        raise ValueError(
                            f"Partition column '{col}' must be present in source for merge"
                        )
                    eq = pc.equal(joined.column(col), joined.column(src_name))
                    neq = pc.invert(eq)
                    violations = pc.and_(match_mask, pc.fill_null(neq, True))
                    if pc.any(violations).as_py():
                        raise ValueError(
                            "Cannot merge: partition column values cannot change for existing keys"
                        )

            out_arrays = []
            out_names = []
            for col in output_columns:
                if col in key_cols:
                    out_arrays.append(joined.column(col))
                    out_names.append(col)
                    continue

                src_name = f"{col}__src"
                if src_name in joined.column_names:
                    out_arrays.append(
                        pc.if_else(
                            match_mask, joined.column(src_name), joined.column(col)
                        )
                    )
                else:
                    out_arrays.append(joined.column(col))
                out_names.append(col)

            updated_table = pa_mod.table(out_arrays, names=out_names)

            staging_file = f"{staging_dir}/{uuid.uuid4().hex[:16]}.parquet"
            pq.write_table(
                updated_table,
                staging_file,
                filesystem=fs,
                compression=compression,
                row_group_size=row_group_size,
            )

            size_bytes = None
            try:
                size_bytes = int(fs.size(staging_file))
            except (OSError, IOError, PermissionError):
                size_bytes = None

            file_manager.atomic_replace_files(
                [staging_file], [file_path], filesystem=fs
            )

            rewritten_files.append(file_path)
            rewritten_meta.append(
                MergeFileMetadata(
                    path=file_path,
                    row_count=updated_table.num_rows,
                    operation="rewritten",
                    size_bytes=size_bytes,
                )
            )
    finally:
        file_manager.cleanup_staging_files(filesystem=fs)

    inserted_files: list[str] = []
    inserted_meta: list[MergeFileMetadata] = []
    inserted_rows = 0

    if strategy == "upsert" and inserted_key_set:
        insert_table = self._select_rows_by_keys(
            source_table,
            key_cols,
            set(inserted_key_set),
        )
        inserted_rows = insert_table.num_rows
        write_res = self.write_dataset(
            insert_table,
            path,
            mode="append",
            partition_by=partition_cols or None,
            compression=compression,
            max_rows_per_file=max_rows_per_file,
            row_group_size=row_group_size,
        )
        inserted_files = [m.path for m in write_res.files]
        inserted_meta = [
            MergeFileMetadata(
                path=m.path,
                row_count=m.row_count,
                operation="inserted",
                size_bytes=m.size_bytes,
            )
            for m in write_res.files
        ]

    updated_rows = len(matched_keys)

    files_meta = (
        rewritten_meta
        + inserted_meta
        + [
            MergeFileMetadata(path=f, row_count=0, operation="preserved")
            for f in preserved_files
        ]
    )

    return MergeResult(
        strategy=strategy,
        source_count=source_table.num_rows,
        target_count_before=target_count_before,
        target_count_after=target_count_before + inserted_rows,
        inserted=inserted_rows,
        updated=updated_rows if strategy != "insert" else 0,
        deleted=0,
        files=files_meta,
        rewritten_files=rewritten_files,
        inserted_files=inserted_files,
        preserved_files=preserved_files,
    )
fsspeckit.datasets.duckdb.DuckDBDatasetIO.read_parquet
read_parquet(
    path: str,
    columns: list[str] | None = None,
    filters: Any | None = None,
    use_threads: bool = True,
) -> Table

Read parquet file(s) using DuckDB.

Parameters:

Name Type Description Default
path str

Path to parquet file or directory

required
columns list[str] | None

Optional list of columns to read

None
filters Any | None

Optional SQL WHERE clause string for DuckDB (e.g., "column > 5 AND other = 'value'")

None
use_threads bool

Whether to use parallel reading (DuckDB ignores this)

True

Returns:

Type Description
Table

PyArrow table containing the data

Raises:

Type Description
TypeError

If filters is not None and not a string

Example
1
2
3
4
5
6
from fsspeckit.datasets.duckdb.connection import create_duckdb_connection
from fsspeckit.datasets.duckdb.dataset import DuckDBDatasetIO

conn = create_duckdb_connection()
io = DuckDBDatasetIO(conn)
table = io.read_parquet("/path/to/file.parquet", filters="id > 100")
Source code in src/fsspeckit/datasets/duckdb/dataset.py
def read_parquet(
    self,
    path: str,
    columns: list[str] | None = None,
    filters: Any | None = None,
    use_threads: bool = True,
) -> pa.Table:
    """Read parquet file(s) using DuckDB.

    Args:
        path: Path to parquet file or directory
        columns: Optional list of columns to read
        filters: Optional SQL WHERE clause string for DuckDB (e.g., "column > 5 AND other = 'value'")
        use_threads: Whether to use parallel reading (DuckDB ignores this)

    Returns:
        PyArrow table containing the data

    Raises:
        TypeError: If filters is not None and not a string

    Example:
        ```python
        from fsspeckit.datasets.duckdb.connection import create_duckdb_connection
        from fsspeckit.datasets.duckdb.dataset import DuckDBDatasetIO

        conn = create_duckdb_connection()
        io = DuckDBDatasetIO(conn)
        table = io.read_parquet("/path/to/file.parquet", filters="id > 100")
        ```
    """
    validate_path(path)

    if filters is not None and not isinstance(filters, str):
        raise TypeError(
            "DuckDB filters must be a SQL WHERE clause string. "
            "Received type: {type(filters).__name__}. "
            "Example: filters='column > 5 AND other = \"value\"'"
        )

    conn = self._connection.connection

    # Build the query
    query = "SELECT * FROM parquet_scan(?)"

    params = [path]

    if columns:
        # Escape column names and build select list
        quoted_cols = [f'"{col}"' for col in columns]
        select_list = ", ".join(quoted_cols)
        query = f"SELECT {select_list} FROM parquet_scan(?)"

    if filters:
        query += f" WHERE {filters}"

    # DuckDB ignores use_threads parameter, but we accept it for interface compatibility
    _ = use_threads

    try:
        # Execute query
        result = conn.execute(query, params).fetch_arrow_table()

        return result

    except (
        _DUCKDB_EXCEPTIONS.get("IOException"),
        _DUCKDB_EXCEPTIONS.get("InvalidInputException"),
        _DUCKDB_EXCEPTIONS.get("ParserException"),
    ) as e:
        raise RuntimeError(
            f"Failed to read parquet from {path}: {safe_format_error(e)}"
        ) from e
fsspeckit.datasets.duckdb.DuckDBDatasetIO.write_dataset
write_dataset(
    data: Table | list[Table],
    path: str,
    *,
    mode: Literal["append", "overwrite"] = "append",
    basename_template: str | None = None,
    schema: Schema | None = None,
    partition_by: str | list[str] | None = None,
    compression: str | None = "snappy",
    max_rows_per_file: int | None = 5000000,
    row_group_size: int | None = 500000,
) -> "WriteDatasetResult"

Write a parquet dataset and return per-file metadata.

Source code in src/fsspeckit/datasets/duckdb/dataset.py
def write_dataset(
    self,
    data: pa.Table | list[pa.Table],
    path: str,
    *,
    mode: Literal["append", "overwrite"] = "append",
    basename_template: str | None = None,
    schema: pa.Schema | None = None,
    partition_by: str | list[str] | None = None,
    compression: str | None = "snappy",
    max_rows_per_file: int | None = 5_000_000,
    row_group_size: int | None = 500_000,
) -> "WriteDatasetResult":
    """Write a parquet dataset and return per-file metadata."""
    import uuid

    from fsspeckit.common.security import validate_compression_codec, validate_path
    from fsspeckit.core.incremental import IncrementalFileManager
    from fsspeckit.datasets.write_result import (
        FileWriteMetadata,
        WriteDatasetResult,
    )

    validate_path(path)
    validate_compression_codec(compression)

    self._validate_write_mode(mode)
    row_group_size = self._validate_write_parameters(
        max_rows_per_file,
        row_group_size,
    )

    table = self._combine_tables(data)
    if schema is not None:
        from fsspeckit.datasets.schema import cast_schema

        table = cast_schema(table, schema)

    partition_cols = self._validate_partition_columns(
        partition_by,
        table.column_names,
    )

    if basename_template is None:
        basename_template = "part-{i}.parquet"

    if mode == "append" and basename_template == "part-{i}.parquet":
        unique_id = uuid.uuid4().hex[:16]
        basename_template = f"part-{unique_id}-{{i}}.parquet"

    def _format_filename(index: int) -> str:
        if "{i}" in basename_template:
            return basename_template.format(i=index)
        if basename_template.endswith(".parquet"):
            stem = basename_template[:-8]
            return f"{stem}-{uuid.uuid4().hex[:16]}.parquet"
        return f"{basename_template}-{uuid.uuid4().hex[:16]}"

    fs = self._connection.filesystem
    fs.mkdirs(path, exist_ok=True)

    if mode == "overwrite":
        self._clear_dataset_parquet_only(path)

    file_manager = IncrementalFileManager()
    staging_dir = file_manager.create_staging_directory(path, filesystem=fs)

    moved_files: list[str] = []
    try:
        self._write_to_path(
            data=table,
            path=staging_dir,
            compression=compression,
            max_rows_per_file=max_rows_per_file,
            row_group_size=row_group_size,
            mode="overwrite",
            partition_by=partition_cols or None,
        )

        staging_files = [
            f
            for f in fs.find(staging_dir, withdirs=False)
            if f.endswith(".parquet")
        ]
        staging_prefix = staging_dir.rstrip("/") + "/"

        for index, staging_file in enumerate(staging_files):
            staging_file_path = fs._strip_protocol(staging_file)
            staging_prefix_path = fs._strip_protocol(staging_prefix)

            if os.path.isabs(staging_file_path) and not os.path.isabs(
                staging_prefix_path
            ):
                staging_prefix_path = os.path.abspath(staging_prefix_path)

            if staging_file_path.startswith(staging_prefix_path):
                relative = staging_file_path[len(staging_prefix_path) :].lstrip("/")
            else:
                relative = os.path.relpath(staging_file_path, staging_prefix_path)

            relative_path = Path(relative)
            partition_dir = relative_path.parent.as_posix()

            if partition_dir not in ("", "."):
                partition_parts = relative_path.parent.parts
                if (
                    partition_cols
                    and len(partition_parts) == len(partition_cols)
                    and not any("=" in part for part in partition_parts)
                ):
                    # Normalize value-only directories to Hive-style col=value
                    partition_dir = "/".join(
                        f"{col}={val}"
                        for col, val in zip(partition_cols, partition_parts)
                    )

            target_dir = (
                path if partition_dir in ("", ".") else f"{path}/{partition_dir}"
            )
            fs.mkdirs(target_dir, exist_ok=True)
            filename = _format_filename(index)
            target_file = f"{target_dir}/{filename}"
            fs.move(staging_file, target_file)
            moved_files.append(target_file)
    finally:
        file_manager.cleanup_staging_files(filesystem=fs)

    files: list[FileWriteMetadata] = []
    for f in moved_files:
        row_count = int(self._get_file_row_count(f))
        size_bytes = None
        try:
            size_bytes = int(fs.size(f))
        except (OSError, IOError, PermissionError) as e:
            logger.warning(
                "Failed to retrieve file size",
                path=f,
                error=str(e),
                operation="write_dataset",
            )
            size_bytes = None
        except (TypeError, ValueError) as e:
            logger.warning(
                "Invalid file size value",
                path=f,
                error=str(e),
                operation="write_dataset",
            )
            size_bytes = None

        files.append(
            FileWriteMetadata(path=f, row_count=row_count, size_bytes=size_bytes)
        )

    return WriteDatasetResult(
        files=files,
        total_rows=sum(f.row_count for f in files),
        mode=mode,
        backend="duckdb",
    )
fsspeckit.datasets.duckdb.DuckDBDatasetIO.write_parquet
write_parquet(
    data: Table | list[Table],
    path: str,
    compression: str | None = "snappy",
    row_group_size: int | None = None,
    use_threads: bool = False,
) -> None

Write parquet file using DuckDB.

Parameters:

Name Type Description Default
data Table | list[Table]

PyArrow table or list of tables to write

required
path str

Output file path

required
compression str | None

Compression codec to use

'snappy'
row_group_size int | None

Rows per row group

None
use_threads bool

Whether to use parallel writing

False
Example
1
2
3
4
5
6
7
8
import pyarrow as pa
from fsspeckit.datasets.duckdb.connection import create_duckdb_connection
from fsspeckit.datasets.duckdb.dataset import DuckDBDatasetIO

table = pa.table({'a': [1, 2, 3], 'b': ['x', 'y', 'z']})
conn = create_duckdb_connection()
io = DuckDBDatasetIO(conn)
io.write_parquet(table, "/tmp/data.parquet")
Source code in src/fsspeckit/datasets/duckdb/dataset.py
def write_parquet(
    self,
    data: pa.Table | list[pa.Table],
    path: str,
    compression: str | None = "snappy",
    row_group_size: int | None = None,
    use_threads: bool = False,
) -> None:
    """Write parquet file using DuckDB.

    Args:
        data: PyArrow table or list of tables to write
        path: Output file path
        compression: Compression codec to use
        row_group_size: Rows per row group
        use_threads: Whether to use parallel writing

    Example:
        ```python
        import pyarrow as pa
        from fsspeckit.datasets.duckdb.connection import create_duckdb_connection
        from fsspeckit.datasets.duckdb.dataset import DuckDBDatasetIO

        table = pa.table({'a': [1, 2, 3], 'b': ['x', 'y', 'z']})
        conn = create_duckdb_connection()
        io = DuckDBDatasetIO(conn)
        io.write_parquet(table, "/tmp/data.parquet")
        ```
    """
    validate_path(path)
    compression_final = compression or "snappy"
    validate_compression_codec(compression_final)

    fs = self._connection.filesystem
    parent = str(Path(path).parent)
    if parent and parent not in (".", "/"):
        fs.mkdirs(parent, exist_ok=True)

    conn = self._connection.connection
    table = self._combine_tables(data)

    # Register the data as a temporary table
    f"temp_{uuid.uuid4().hex[:16]}"
    conn.register("data_table", table)

    try:
        # Build the COPY command
        copy_query = "COPY data_table TO ?"

        params = [path]

        options: list[str] = []
        if compression_final:
            options.append(f"COMPRESSION {compression_final}")
        if row_group_size:
            options.append(f"ROW_GROUP_SIZE {row_group_size}")
        if options:
            copy_query += " (" + ", ".join(options) + ")"

        # Execute the copy
        if use_threads:
            conn.execute(copy_query, params)
        else:
            conn.execute(copy_query, params)

    finally:
        # Clean up temporary table
        _unregister_duckdb_table_safely(conn, "data_table")

fsspeckit.datasets.duckdb.MergeStrategy

Bases: Enum

Supported merge strategies with consistent semantics across backends.

Attributes
fsspeckit.datasets.duckdb.MergeStrategy.DEDUPLICATE class-attribute instance-attribute
DEDUPLICATE = 'deduplicate'

Remove duplicates from source, then upsert.

fsspeckit.datasets.duckdb.MergeStrategy.FULL_MERGE class-attribute instance-attribute
FULL_MERGE = 'full_merge'

Insert, update, and delete (full sync with source).

fsspeckit.datasets.duckdb.MergeStrategy.INSERT class-attribute instance-attribute
INSERT = 'insert'

Insert only new records, ignore existing records.

fsspeckit.datasets.duckdb.MergeStrategy.UPDATE class-attribute instance-attribute
UPDATE = 'update'

Update only existing records, ignore new records.

fsspeckit.datasets.duckdb.MergeStrategy.UPSERT class-attribute instance-attribute
UPSERT = 'upsert'

Insert new records, update existing records.

Functions

fsspeckit.datasets.duckdb.create_duckdb_connection

create_duckdb_connection(
    filesystem: AbstractFileSystem | None = None,
) -> DuckDBConnection

Create a DuckDB connection manager.

Parameters:

Name Type Description Default
filesystem AbstractFileSystem | None

fsspec filesystem to use

None

Returns:

Type Description
DuckDBConnection

DuckDB connection manager

Source code in src/fsspeckit/datasets/duckdb/connection.py
def create_duckdb_connection(
    filesystem: AbstractFileSystem | None = None,
) -> DuckDBConnection:
    """Create a DuckDB connection manager.

    Args:
        filesystem: fsspec filesystem to use

    Returns:
        DuckDB connection manager
    """
    return DuckDBConnection(filesystem=filesystem)

fsspeckit.datasets.duckdb.collect_dataset_stats_duckdb

collect_dataset_stats_duckdb(
    path: str,
    filesystem: AbstractFileSystem | None = None,
    partition_filter: list[str] | None = None,
) -> dict[str, Any]

Collect file-level statistics for a parquet dataset using shared core logic.

This function delegates to the shared fsspeckit.core.maintenance.collect_dataset_stats function, ensuring consistent dataset discovery and statistics across both DuckDB and PyArrow backends.

The helper walks the given dataset directory on the provided filesystem, discovers parquet files (recursively), and returns basic statistics:

  • Per-file path, size in bytes, and number of rows
  • Aggregated total bytes and total rows

The function is intentionally streaming/metadata-driven and never materializes the full dataset as a single table.

Parameters:

Name Type Description Default
path str

Root directory of the parquet dataset.

required
filesystem AbstractFileSystem | None

Optional fsspec filesystem. If omitted, a local "file" filesystem is used.

None
partition_filter list[str] | None

Optional list of partition prefix filters (e.g. ["date=2025-11-04"]). Only files whose path relative to path starts with one of these prefixes are included.

None

Returns:

Type Description
dict[str, Any]

Dict with keys:

dict[str, Any]
  • files: list of {"path", "size_bytes", "num_rows"} dicts
dict[str, Any]
  • total_bytes: sum of file sizes
dict[str, Any]
  • total_rows: sum of row counts

Raises:

Type Description
FileNotFoundError

If the path does not exist or no parquet files match the optional partition filter.

Note

This is a thin wrapper around the shared core function. See :func:fsspeckit.core.maintenance.collect_dataset_stats for the authoritative implementation.

Source code in src/fsspeckit/datasets/duckdb/dataset.py
def collect_dataset_stats_duckdb(
    path: str,
    filesystem: AbstractFileSystem | None = None,
    partition_filter: list[str] | None = None,
) -> dict[str, Any]:
    """Collect file-level statistics for a parquet dataset using shared core logic.

    This function delegates to the shared ``fsspeckit.core.maintenance.collect_dataset_stats``
    function, ensuring consistent dataset discovery and statistics across both DuckDB
    and PyArrow backends.

    The helper walks the given dataset directory on the provided filesystem,
    discovers parquet files (recursively), and returns basic statistics:

    - Per-file path, size in bytes, and number of rows
    - Aggregated total bytes and total rows

    The function is intentionally streaming/metadata-driven and never
    materializes the full dataset as a single table.

    Args:
        path: Root directory of the parquet dataset.
        filesystem: Optional fsspec filesystem. If omitted, a local "file"
            filesystem is used.
        partition_filter: Optional list of partition prefix filters
            (e.g. ["date=2025-11-04"]). Only files whose path relative to
            ``path`` starts with one of these prefixes are included.

    Returns:
        Dict with keys:

        - ``files``: list of ``{"path", "size_bytes", "num_rows"}`` dicts
        - ``total_bytes``: sum of file sizes
        - ``total_rows``: sum of row counts

    Raises:
        FileNotFoundError: If the path does not exist or no parquet files
            match the optional partition filter.

    Note:
        This is a thin wrapper around the shared core function. See
        :func:`fsspeckit.core.maintenance.collect_dataset_stats` for the
        authoritative implementation.
    """
    from fsspeckit.core.maintenance import collect_dataset_stats

    return collect_dataset_stats(
        path=path,
        filesystem=filesystem,
        partition_filter=partition_filter,
    )