Skip to content
Open
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
27 changes: 21 additions & 6 deletions examples/csv-read-options.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,15 +17,30 @@

"""Example demonstrating CsvReadOptions usage."""

import gzip
import shutil
import tempfile
from pathlib import Path

from datafusion import CsvReadOptions, SessionContext

# This example is self-contained: it writes its own small CSV file, plus a
# gzipped copy, into a temporary directory instead of expecting a
# ``data.csv`` file that is not part of this repository.
tmp_dir = Path(tempfile.mkdtemp())
csv_path = tmp_dir / "data.csv"
csv_path.write_text("a,b,c\n1,2,3\n4,5,6\n7,8,9\n")
gz_path = tmp_dir / "data.csv.gz"
with csv_path.open("rb") as f_in, gzip.open(gz_path, "wb") as f_out:
shutil.copyfileobj(f_in, f_out)

# Create a SessionContext
ctx = SessionContext()

# Example 1: Using CsvReadOptions with default values
print("Example 1: Default CsvReadOptions")
options = CsvReadOptions()
df = ctx.read_csv("data.csv", options=options)
df = ctx.read_csv(str(csv_path), options=options)

# Example 2: Using CsvReadOptions with custom parameters
print("\nExample 2: Custom CsvReadOptions")
Expand All @@ -36,7 +51,7 @@
schema_infer_max_records=1000,
file_extension=".csv",
)
df = ctx.read_csv("data.csv", options=options)
df = ctx.read_csv(str(csv_path), options=options)

# Example 3: Using the builder pattern (recommended for readability)
print("\nExample 3: Builder pattern")
Expand All @@ -49,7 +64,7 @@
.with_truncated_rows(False) # noqa: FBT003
.with_newlines_in_values(True) # noqa: FBT003
)
df = ctx.read_csv("data.csv", options=options)
df = ctx.read_csv(str(csv_path), options=options)

# Example 4: Advanced options
print("\nExample 4: Advanced options")
Expand All @@ -64,18 +79,18 @@
.with_file_compression_type("gzip") # Read gzipped CSV
.with_file_extension(".gz")
)
df = ctx.read_csv("data.csv.gz", options=options)
df = ctx.read_csv(str(gz_path), options=options)

# Example 5: Register CSV table with options
print("\nExample 5: Register CSV table")
options = CsvReadOptions().with_has_header(True).with_delimiter(",") # noqa: FBT003
ctx.register_csv("my_table", "data.csv", options=options)
ctx.register_csv("my_table", str(csv_path), options=options)
df = ctx.sql("SELECT * FROM my_table")

# Example 6: Backward compatibility (without options)
print("\nExample 6: Backward compatibility")
# Still works the old way!
df = ctx.read_csv("data.csv", has_header=True, delimiter=",")
df = ctx.read_csv(str(csv_path), has_header=True, delimiter=",")

print("\nAll examples completed!")
print("\nFor all available options, see the CsvReadOptions documentation:")
Expand Down
5 changes: 5 additions & 0 deletions examples/export.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,19 +34,24 @@
# export to pandas dataframe
pandas_df = df.to_pandas()
assert pandas_df.shape == (3, 2)
print(pandas_df)

# export to PyArrow table
arrow_table = df.to_arrow_table()
assert arrow_table.shape == (3, 2)
print(arrow_table)

# export to Polars dataframe
polars_df = df.to_polars()
assert polars_df.shape == (3, 2)
print(polars_df)

# export to Python list of rows
pylist = df.to_pylist()
assert pylist == [{"a": 1, "b": 4}, {"a": 2, "b": 5}, {"a": 3, "b": 6}]
print(pylist)

# export to Python dictionary of columns
pydict = df.to_pydict()
assert pydict == {"a": [1, 2, 3], "b": [4, 5, 6]}
print(pydict)
4 changes: 4 additions & 0 deletions examples/import.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,18 +40,22 @@
# Create a datafusion DataFrame from a Python list of rows
df = ctx.from_pylist([{"a": 1, "b": 4}, {"a": 2, "b": 5}, {"a": 3, "b": 6}])
assert type(df) is datafusion.DataFrame
df.show()

# Convert pandas DataFrame to datafusion DataFrame
pandas_df = pd.DataFrame({"a": [1, 2, 3], "b": [4, 5, 6]})
df = ctx.from_pandas(pandas_df)
assert type(df) is datafusion.DataFrame
df.show()

# Convert polars DataFrame to datafusion DataFrame
polars_df = pl.DataFrame({"a": [1, 2, 3], "b": [4, 5, 6]})
df = ctx.from_polars(polars_df)
assert type(df) is datafusion.DataFrame
df.show()

# Convert Arrow Table to datafusion DataFrame
arrow_table = pa.Table.from_pydict({"a": [1, 2, 3], "b": [4, 5, 6]})
df = ctx.from_arrow(arrow_table)
assert type(df) is datafusion.DataFrame
df.show()
1 change: 1 addition & 0 deletions examples/python-udaf.py
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ def evaluate(self) -> pa.Scalar:
)

df = df.aggregate([], [my_udaf(col("a"))])
df.show()

result = df.collect()[0]

Expand Down
1 change: 1 addition & 0 deletions examples/python-udf.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ def is_null(array: pa.Array) -> pa.Array:
df = ctx.create_dataframe([[batch]])

df = df.select(is_null_arr(f.col("a")))
df.show()

result = df.collect()[0]

Expand Down
2 changes: 2 additions & 0 deletions examples/query-pyarrow-data.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,8 @@
col("a") - col("b"),
)

df.show()

# execute and collect the first (and only) batch
result = df.collect()[0]

Expand Down
1 change: 1 addition & 0 deletions examples/sql-to-pandas.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@

# convert to Pandas
pandas_df = df.to_pandas()
print(pandas_df)

# create a chart
fig = pandas_df.plot(
Expand Down
1 change: 1 addition & 0 deletions examples/sql-using-python-udaf.py
Original file line number Diff line number Diff line change
Expand Up @@ -82,5 +82,6 @@ def evaluate(self) -> pa.Scalar:
# | 1 | 9 |
# | 3 | 6 |
# +---+--------------+
result_df.show()
assert result_df.to_pydict()["a"] == [1, 3]
assert result_df.to_pydict()["b_aggregated"] == [9, 6]
1 change: 1 addition & 0 deletions examples/sql-using-python-udf.py
Original file line number Diff line number Diff line change
Expand Up @@ -61,4 +61,5 @@ def is_null(array: pa.Array) -> pa.Array:
# | 2 | true |
# | 3 | false |
# +---+-----------+
result_df.show()
assert result_df.to_pydict()["b_is_null"] == [False, True, False]
3 changes: 3 additions & 0 deletions examples/substrait.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,3 +47,6 @@
# Back to Substrait Plan just for demonstration purposes
# type(substrait_plan) -> <class 'datafusion.substrait.plan'>
substrait_plan = ss.Producer.to_substrait_plan(df_logical_plan, ctx)

# Show the logical plan recovered from the Substrait round trip
print(df_logical_plan)