| 文件 | 最后提交记录 | 最后更新时间 |
|---|---|---|
[RUST] Migrate legacy tests (#775) Migration of legacy tests from tests_legacy/ to tests/ * tests_legacy/null_tests/ are now moved into their respective tests in tests/dataframe/*.py * Deleted tests_legacy/expression_operators/* tests, which should now be tested on the Series-level * Deleted tests_legacy/experimental/* tests * Deleted tests_legacy/internal since these are no longer used from Python * Deleted tests_legacy/runners since these are redundant and don't serve much purpose --------- Co-authored-by: Jay Chia <jaychia94@gmail.com@users.noreply.github.com> | 3 年前 | |
[RUST] Migrate legacy tests (#775) Migration of legacy tests from tests_legacy/ to tests/ * tests_legacy/null_tests/ are now moved into their respective tests in tests/dataframe/*.py * Deleted tests_legacy/expression_operators/* tests, which should now be tested on the Series-level * Deleted tests_legacy/experimental/* tests * Deleted tests_legacy/internal since these are no longer used from Python * Deleted tests_legacy/runners since these are redundant and don't serve much purpose --------- Co-authored-by: Jay Chia <jaychia94@gmail.com@users.noreply.github.com> | 3 年前 | |
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> | 9 个月前 | |
feat(expressions): support count(mode='all') without expr (#6358) ## Changes Made This PR adds support for calling daft.functions.count(mode="all") without passing an expression, enabling free-standing row-count aggregations similar to SQL COUNT(*). ### What changed - Updated daft.functions.count in daft/functions/agg.py to accept an optional expr argument. - Added explicit behavior for expr=None: - mode="all" is allowed and maps to count(*) semantics via col("*"). - mode="valid" or mode="null" now raises a clear ValueError when no expression is provided. - Preserved backward compatibility for all existing count(expr, mode=...) and count(lit(...)) usages. ### Tests added In tests/dataframe/test_aggregations.py: - Global aggregation with count(mode="all") and no expression. - Grouped aggregation with count(mode="all") and no expression. - Empty DataFrame behavior (0 result). - Error path validation for no-expression calls with non-all modes. ### Validation run - DAFT_RUNNER=native make test EXTRA_ARGS="-v tests/dataframe/test_aggregations.py" - DAFT_RUNNER=native make test EXTRA_ARGS="-v tests/sql/test_sql.py tests/sql/test_exprs.py" ## Related Issues Closes #5526 Co-authored-by: Lord of Abyss <pancx@chinatelecom.cn> | 6 个月前 | |
refactor(arrow-rs): Migrate approx count distinct (#6038) ## Changes Made Removes the use of arrow2 primitive arrays. | 7 个月前 | |
[BUG] Enable groupby with alias for native executor (#2917) Fixed this already for the current translation path back in https://github.com/Eventual-Inc/Daft/pull/2790, doing the same for swordfish --------- Co-authored-by: Colin Ho <colinho@Colins-MacBook-Pro.local> Co-authored-by: Colin Ho <colinho@Colins-MBP.localdomain> | 1 年前 | |
perf: Call async udfs asynchronously (#5451) ## Changes Made Call async UDFs from swordfish asynchronously using pyo3-async-runtimes. This allows async UDFs to actually run asynchronously without blocking a thread. Consider the following example, where we run 2 async udfs. On main, currently only a total of num_cores across _both_ async udfs can run concurrently, but with this PR, we can run all concurrently as long as they have input ready. Running on my macbook the speed up is ~20s -> ~2s. import daft import asyncio @daft.func async def my_udf(foo: int) -> int: await asyncio.sleep(1) return foo + 1 df = ( daft.from_pydict({"foo": [i for i in range(1000)]}).into_batches(10) .select(my_udf(daft.col("foo"))) .select(my_udf(daft.col("foo"))) .collect() ) ## Implementation Details 1. Convert def call_async_batch to an async def call_async_batch 2. Add an async fn call_async to row wise python udf that uses pyo3_async_runtimes to call async def call_async_batch 3. Create a new async_udf streaming sink. 4. Call the new async fn call_async when evaluating the udf in async_udf operator, spawn the future on a joinset and poll it during calls to execute. Notes: - We have to store a reference to a running asyncio event loop in order to use pyo3-async-runtimes. This works natively on flotilla because we use async Ray actors, but on swordfish we need to create one ourselves, so I opted to create a background event loop on a separate python thread. - I verified that this approach works even on Jupyter which by default runs it's own event loop - The joinset approach will allow the udfs to be driven asynchronously, but their results will only be retrieved during calls to execute / finalize. ## Todos - yield incrementally from streaming sink finalize - consider using orderedjoinset to maintain order - batch async udfs - user defined concurrency control ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> ## Checklist - [ ] Documented in API Docs (if applicable) - [ ] Documented in User Guide (if applicable) - [ ] If adding a new documentation page, doc is added to docs/mkdocs.yml navigation - [ ] Documentation builds and is formatted properly | 10 个月前 | |
[CHORE] Swordfish specific test fixtures (#3164) This PR sets up a few swordfish related test fixtures, specifically: - Parameterize default_morsel_size = [1, None] for dataframe tests that do any into/repartitioning. This is to make sure that the operator parallelism is working. - Setup iteration tests in test_iter.py - Makes the ordering assertions stricter on some tests. E.g. some tests do assert df.sort(col) == expected, but there are other columns in df that may not be sorted, and this won't be enough if morsel_size = 1. This isn't a problem with swordfish but the test, where the sort should actually involve more columns. Big note: There was a problem with pivot not getting applied correctly. This is because a dataframe pivot operation comprises of an agg + the actual pivoting, but previously the pivot was implemented as an intermediate operator, and the results of the agg were getting buffered. In order for the pivot to work it has to receive all values with the same group_by keys. This PR implements simplifies Pivot as a BlockingSink, so all the work is in there. --------- Co-authored-by: Colin Ho <colinho@Colins-MacBook-Pro.local> | 1 年前 | |
feat: support count(1) in dataframe and choose the cheap column (#4977) ## Changes Made <!-- Describe what changes were made and why. Include implementation details if necessary. --> Not count(1) in dataframe is not supported: In [49]: df.count(1).show() --------------------------------------------------------------------------- TypeError Traceback (most recent call last) Cell In[49], line 1 ----> 1 df.count(1).show() [... skipping hidden 1 frame] File /data00/code/tmp2/Daft/daft/dataframe/dataframe.py:3011, in DataFrame.count(self, *cols) 3008 raise ValueError("Cannot call count() with both * and column names") 3010 # Otherwise, perform a column-wise count on the specified columns -> 3011 return self._apply_agg_fn(Expression.count, cols) File /data00/code/tmp2/Daft/daft/dataframe/dataframe.py:2854, in DataFrame._apply_agg_fn(self, fn, cols, group_by) 2852 groupby_name_set = set() if group_by is None else group_by.to_name_set() 2853 cols = tuple(c for c in self.column_names if c not in groupby_name_set) -> 2854 exprs = self._wildcard_inputs_to_expressions(cols) 2855 return self._agg([fn(c) for c in exprs], group_by) File /data00/code/tmp2/Daft/daft/dataframe/dataframe.py:1596, in DataFrame._wildcard_inputs_to_expressions(self, columns) 1594 """Handles wildcard argument column inputs.""" 1595 column_input: Iterable[ColumnInputType] = columns[0] if len(columns) == 1 else columns # type: ignore -> 1596 return column_inputs_to_expressions(column_input) File /data00/code/tmp2/Daft/daft/utils.py:126, in column_inputs_to_expressions(columns) 123 from daft.expressions import col 125 column_iter: Iterable[ColumnInputType] = [columns] if is_column_input(columns) else columns # type: ignore --> 126 return [col(c) if isinstance(c, str) else c for c in column_iter] TypeError: 'int' object is not iterable Using this pr: In [5]: df.count(1).show() ╭────────╮ │ count │ │ --- │ │ UInt64 │ ╞════════╡ │ 4 │ ╰────────╯ At the same time, there is also a certain improvement in performance. Before modification, the performance data is as follows -------------------------------------------------------------------------------------------------------- benchmark: 4 tests ------------------------------------------------------------------------------------------------------- Name (time in ms) Min Max Mean StdDev Median IQR Outliers OPS Rounds Iterations ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- test_count_with_col_name[float_id-float] 18.2832 (1.0) 24.8761 (1.07) 20.1458 (1.0) 2.6823 (1.34) 19.3201 (1.0) 2.0831 (1.0) 1;1 49.6381 (1.0) 5 1 test_count_with_col_name[int_id-int] 18.6567 (1.02) 23.2190 (1.0) 21.0213 (1.04) 1.9988 (1.0) 21.3333 (1.10) 3.6173 (1.74) 2;0 47.5707 (0.96) 5 1 test_count_with_col_name[string_content-string] 314.5036 (17.20) 633.3133 (27.28) 434.0078 (21.54) 120.1076 (60.09) 414.9809 (21.48) 120.6531 (57.92) 1;0 2.3041 (0.05) 5 1 test_count_with_col_name[binary_content-binary] 1,836.5554 (100.45) 1,889.4419 (81.37) 1,864.8847 (92.57) 22.0420 (11.03) 1,857.1332 (96.12) 34.7329 (16.67) 2;0 0.5362 (0.01) 5 1 ----- Use this pr, the performance indicators are as follows: ------------------------------------------------------------------------------------------------ benchmark: 4 tests ----------------------------------------------------------------------------------------------- Name (time in ms) Min Max Mean StdDev Median IQR Outliers OPS Rounds Iterations ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- test_count_with_col_name[int_id-int] 18.3580 (1.0) 25.2071 (1.00) 21.3138 (1.0) 2.7985 (1.68) 20.9401 (1.0) 4.5452 (2.61) 2;0 46.9180 (1.0) 5 1 test_count_with_col_name[float_id-float] 20.5159 (1.12) 25.1621 (1.0) 22.8351 (1.07) 1.6660 (1.0) 22.9081 (1.09) 1.7429 (1.0) 2;0 43.7923 (0.93) 5 1 test_count_with_col_name[binary_content-binary] 21.8774 (1.19) 27.1461 (1.08) 23.9783 (1.13) 2.1795 (1.31) 24.3946 (1.16) 3.2378 (1.86) 1;0 41.7043 (0.89) 5 1 test_count_with_col_name[string_content-string] 22.6195 (1.23) 30.9669 (1.23) 25.7760 (1.21) 3.5412 (2.13) 23.9304 (1.14) 5.4121 (3.11) 1;0 38.7958 (0.83) 5 1 ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- The benchmark test cases are as follows: from __future__ import annotations import pytest import daft PATH = "/tmp/test_count_1" # Consolidated test configuration with built-in IDs TEST_CASES = [ pytest.param("int_id", "int", id="first_col_int"), pytest.param("float_id", "float", id="first_col_float"), pytest.param("string_content", "string", id="first_col_string"), pytest.param("binary_content", "binary", id="first_col_binary"), ] def generate_data(): # Adjust column order to ensure the target column is first data = { "int_id": [1] * 1000, "float_id": [0.001] * 1000, "string_content": ["a"* 100000] * 1000, "binary_content": [b"a" * 1000000] * 1000, } for case in TEST_CASES: param, _ = case.values col_order = [param] + ["int_id", "float_id", "string_content", "binary_content"] col_order = list(dict.fromkeys(col_order)) df = daft.from_pydict({k: data[k] for k in col_order}) path = f"{PATH}_{param}" df.write_parquet(path, write_mode = "overwrite") generate_data() @pytest.mark.parametrize("col_name, _", [case.values[:2] for case in TEST_CASES]) def test_count_with_col_name(benchmark, col_name, _): """Benchmark count(*) with different first column layouts.""" def operation(): path = f"{PATH}_{col_name}" df = daft.read_parquet(path) return df.count().to_pydict() result = benchmark.pedantic(operation, rounds=5, warmup_rounds=1) assert result["count"][0] == 1000 (Showing first 1 of 1 rows) ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> ## Checklist - [ ] Documented in API Docs (if applicable) - [ ] Documented in User Guide (if applicable) - [ ] If adding a new documentation page, doc is added to docs/mkdocs.yml navigation - [ ] Documentation builds and is formatted properly (tag @/ccmao1130 for docs review) --------- Co-authored-by: Colin Ho <colin.ho99@gmail.com> Co-authored-by: EC2 Default User <ec2-user@ip-172-31-22-159.us-west-2.compute.internal> | 1 年前 | |
chore: Upgrade Ruff ruleset to 3.9 and add from __future__ import annotations (#4393) | 1 年前 | |
fix: fix column missing issue for jsonl record schema when reading json/jsonl files (#5681) ## Issue daft.read_json('/mnt/test/vedio.audio-00002.jsonl').where( "sound is null" ).count_rows() ray.exceptions.RayTaskError(DaftCoreException): ray::PhysicalScan->UnGroupedAggregate->Project() (pid=xx, ip=xx.xx.xx.xx, actor_id=xxxx, repr=<daft.runners.flotilla.RaySwordfishActor object at xxxxxxx>) daft.exceptions.DaftCoreException: DaftError::FieldNotFound Column "sound" not found in schema: [Field { name: "video_meta", dtype: Struct([Field { name: "target_fps", dtype: Int64, metadata: {} }]), metadata: {} }] ## Changes Made When reading large JSONL files with automatic schema inference enabled, only the first 1M are used by default to infer the schema. In cases where the data is very sparse, some columns may be missing from the inferred schema. If users later attempt to access such missing columns, an error occurs. The purpose of this change is to eliminate the need for users to manually provide a large and complete schema in sparse-column scenarios. By enhancing the semantics of schema hints to support an "add-missing" behavior, the system can automatically add any missing columns to the final schema when they are accessed. This reduces complexity and improves ease of use. <!-- Describe what changes were made and why. Include implementation details if necessary. --> ## Related Issues https://github.com/Eventual-Inc/Daft/issues/5647 <!-- Link to related GitHub issues, e.g., "Closes #123" --> ## How to test new UTs. | 6 个月前 | |
chore(deps): drop pyarrow 8.0.0 support, bump minimum to >= 15.0.0 (#6378) ## Changes Made Bump minimum PyArrow version from >= 8.0.0 to >= 15.0.0 (44 files, -450 lines). - Update version constraints in pyproject.toml and CI test matrix - Remove _FixSliceOffsets workaround (pyarrow < 12.0.0 struct slice offset bug, fixed upstream) - Remove pyarrow_supports_fixed_shape_tensor() and all conditional branches - Remove obsolete version checks (< 12.0.1, < 13.0.0, >= 9.0.0) and try/except imports - Clean up ~50 pytest.mark.skipif markers and unused PYARROW_GE_* constants in test files Note: _FixEmptyStructArrays is intentionally kept — Daft internally cannot handle empty StructArrays, not just an arrow2 FFI issue. ## Related Issues Closes #6347 | 6 个月前 | |
chore: Upgrade Ruff ruleset to 3.9 and add from __future__ import annotations (#4393) | 1 年前 | |
fix: into_batches should not allow downstream shuffle elision (#6170) ## Changes Made IntoBatchesNode was preserving its child's clustering spec, which could claim num_partitions == 1 even though into_batches could splits data into multiple batches. This caused two cascading failures: 1. distinct saw num_partitions == 1 in needs_hash_repartition and skipped the hash shuffle, so deduplication only happened within each batch and cross-batch duplicates survived. 2. The subsequent count_rows() aggregation also saw num_partitions == 1 in gen_gather_node and skipped the gather step, so each batch produced its own independent count. Python's count_rows() then returned to_pydict()["count"][0] which is the distinct count from just the first batch rather than the global total. We can avoid this by setting the clustering spec to Unknown with num_partitions >= 2. This ensures downstream operators correctly recognize that data is spread across multiple partitions and perform the necessary shuffles/gathers. It's best-effort since there can be more than 2 partitions, but we don't know at planning time how many partitions there will be. ## Related Issues Closes #6161 | 7 个月前 | |
fix(dataframe): explain plan does not make sense when using a write_ op (#6564) ## Problem After calling write_*, explain(show_all=True) collapses to an in-memory source and loses upstream plan context (e.g. filters before write). ## Root Cause write_parquet / write_csv / write_json build the returned DataFrame from write_df._builder after collect(). With cached results, _builder resolves to an in-memory scan builder. ## Solution Build the returned DataFrame from write_df._get_current_builder() so the original logical plan is preserved, while keeping result cache/preview/metadata wiring unchanged. ## Tests - DAFT_RUNNER=native make test EXTRA_ARGS="-v tests/dataframe/test_explain.py::test_explain_after_write_preserves_upstream_plan" ## Impact Fixes explainability for write chains by preserving upstream logical plan visibility without changing write output behavior. --------- Co-authored-by: Lord of Abyss <pancx@chinatelecom.cn> | 5 个月前 | |
fix: add ignore_empty_and_null parameter to .explode() (#6047) Add an optional ignore_empty parameter to explode() that allows users to drop rows with empty lists or null values instead of producing null rows. - Default behavior (ignore_empty=False) is unchanged: empty lists and nulls produce a single null row - When ignore_empty=True: empty lists and nulls are dropped entirely, matching PyArrow/DuckDB behavior ## Changes Made - Added ignore_empty parameter to Python API (DataFrame.explode, Expression.explode, daft.functions.explode) - Propagated the parameter through logical plan, physical plan, and execution layers - Updated Rust core implementation in daft-functions-list and daft-recordbatch - Added 11 test cases covering basic usage and edge cases (all empty lists, all nulls, FixedSizeList, nested lists, etc.) ## Related Issues #5752 --------- Co-authored-by: Srinivas Lade <srinulade1@gmail.com> Co-authored-by: desmondcheongzx <desmondcheongzx@gmail.com> Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com> | 7 个月前 | |
| 7 个月前 | ||
feat: f"" string style string concatenation (#4500) ## Changes Made Implemented f"" string style string concatenation (see linked issue for more detail) and added some tests. Probs want to also add this to the docs at some point? ## Related Issues https://github.com/Eventual-Inc/Daft/issues/4001 ## Checklist - [ ] Documented in API Docs (if applicable) - [ ] Documented in User Guide (if applicable) - [ ] If adding a new documentation page, doc is added to docs/mkdocs.yml navigation - [ ] Documentation builds and is formatted properly (tag @/ccmao1130 for docs review) | 1 年前 | |
refactor!: split column expression into unresolved and resolved types (#3804) BREAKING CHANGE: Python struct get syntactic sugar is no longer allowed (ex: col("a.b") -> col("a").struct.get("b"). I did not think it was worth continuing to maintain what is essentially a small compiler for an undocumented and somewhat unclear feature. I would have to make some significant code changes to get it to work with unresolved columns, and I figured I would just rip it out. In the future I'd like to support something like this instead: col("a")["b"]. --- This PR is the first step to non-equality joins in Daft. We introduce an UnresolvedColumn type so that we can resolve it in the logical plan builder into a ResolvedColumn::Basic, ResolvedColumn::JoinSide (for joins), or ResolvedColumn::OuterRef (for subqueries). In addition, we introduce a SubqueryAlias logical op that is used to determine the sources for UnresolvedColumns with plan_ref set. Along with this PR is a refactor of our SQL planner. The SubqueryAlias op is now used to encode naming and scoping information, so the Relation struct is no longer necessary. Additionally, I discovered that we can significantly simplify SQLPlanner::plan_non_agg_query and SQLPlanner::plan_aggregate_query. I believe the prior complexity came due to the fact that in many clauses (GROUP BY, ORDER BY, etc), you can reference both columns in the input tables as well as aliases in the select clause. However, the planner's bound_columns already solves this issue. This may yield slightly less performant plans due to redundant expressions, but that should be handled by the optimizer instead of the SQL planner. Moreover, I was able to remove a sort from plan_aggregate_query, which may yield both performance and correctness improvements, since the prior code assumed you can split a sort up into two, however that may lead to different results due to sort orders (ORDER BY a, b != ORDER BY b, a) as well as whether our sort impl is stable or unstable. | 1 年前 | |
docs: categorize functions and make each function its own page (#4838) ## Changes Made Change Daft functions docs to categorize them, as I wanted to do with expressions before. Now, the functions are split into different files under daft/functions/, and if you add a function, it will automatically be added to the functions page in the category corresponding to the file it's in! Screenshots: <img width="1584" height="1106" alt="image" src="https://github.com/user-attachments/assets/c5f1d309-1c57-4571-97b6-d303328bfb5a" /> <img width="1584" height="1104" alt="image" src="https://github.com/user-attachments/assets/4edc95bc-f2a1-4d09-97cb-b15eff960a02" /> ## Related Issues #4737 #4824 ## Checklist - [x] Documented in API Docs (if applicable) - [x] Documented in User Guide (if applicable) - [x] If adding a new documentation page, doc is added to docs/mkdocs.yml navigation - [x] Documentation builds and is formatted properly (tag @/ccmao1130 for docs review) | 1 年前 | |
feat(swordfish): plan caching (#6278) ## Summary Introduces fingerprint-based plan caching so multiple executions of the same logical plan share a single pipeline, rather than building a new one each time. ### Key changes - ** PipelineMessage enum** — Replaces raw MicroPartition in all pipeline channels with Morsel { input_id, partition } and Flush(input_id), allowing a single pipeline to multiplex data from multiple logical inputs - **NativeExecutor becomes stateful** — Maintains a HashMap<u64, PlanState> keyed by plan fingerprint; repeated calls to run() with the same fingerprint reuse the existing pipeline - **MessageRouter** — Routes pipeline output to per-input_id unbounded channels so each caller gets only its own results - **try_finish() lifecycle API** — Callers signal input completion and collect per-input stats; pipeline is torn down when the last input finishes - **Per-input runtime stats** — RuntimeStatsManager tracks (NodeID, InputId) pairs with merge support for aggregated views - **next_event helper** — Shared tokio::select! loop abstraction used by blocking sinks, streaming sinks, intermediate ops, and joins - **Simplified BlockingSink trait** — Removed BlockingSinkFinalizeOutput::HasMoreOutput (was dead code); finalize now returns Vec<MicroPartition> directly; make_state takes InputId - **Simplified StreamingSink trait** — Removed StreamingSinkOutput::HasMoreOutput variant - **All node types updated** — Sources, intermediate ops, blocking sinks, streaming sinks, joins, and concat all handle Morsel/Flush per-input routing - **Python integration** — native_executor.py and flotilla.py call try_finish() for lifecycle management ### Files changed ~49 files across daft-local-execution (pipeline, sources, sinks, joins, stats), daft-distributed, common-metrics, and Python runners. ## Test plan - [ ] Full test suite with DAFT_RUNNER=native - [ ] Full test suite with DAFT_RUNNER=ray - [ ] Concurrent query execution sharing the same plan fingerprint - [ ] Verify per-input stats are correctly reported and merged - [ ] TPC-H benchmarks to check for performance regressions 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com> Co-authored-by: R. C. Howell <5731503+rchowell@users.noreply.github.com> | 5 个月前 | |
feat: Flotilla into partitions (#4963) ## Changes Made Implements into_partitions for flotilla. First collect all input tasks. Then either coalesces inputs or splits them such that exactly num_partitions partitions are emitted from this operator. ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> ## Checklist - [ ] Documented in API Docs (if applicable) - [ ] Documented in User Guide (if applicable) - [ ] If adding a new documentation page, doc is added to docs/mkdocs.yml navigation - [ ] Documentation builds and is formatted properly (tag @/ccmao1130 for docs review) | 1 年前 | |
chore(tests): migrate internal usages of daft.udf to cls/func (#6348) update tests to remove a bunch of deprecation warnings in CI/tests ## Changes Made <!-- Describe what changes were made and why. Include implementation details if necessary. --> ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> | 6 个月前 | |
fix: handle case where join keys are different for sort-merge multi-partition join (#6243) ## Changes Made the current sort merge join (multi partition) implementation does not correctly handle the case where the join keys in the left and right dataframes are different. this PR fixes this issue by doing the following: - aliasing the right keys when generating the samples for determining boundaries - renames materialized boundaries with right keys when applying the boundaries to create range partition tasks - regression test added to ensure fix works --------- Co-authored-by: gmweaver <gmweaver.usc@gmail.com> | 7 个月前 | |
| 8 个月前 | ||
chore(deps): drop pyarrow 8.0.0 support, bump minimum to >= 15.0.0 (#6378) ## Changes Made Bump minimum PyArrow version from >= 8.0.0 to >= 15.0.0 (44 files, -450 lines). - Update version constraints in pyproject.toml and CI test matrix - Remove _FixSliceOffsets workaround (pyarrow < 12.0.0 struct slice offset bug, fixed upstream) - Remove pyarrow_supports_fixed_shape_tensor() and all conditional branches - Remove obsolete version checks (< 12.0.1, < 13.0.0, >= 9.0.0) and try/except imports - Clean up ~50 pytest.mark.skipif markers and unused PYARROW_GE_* constants in test files Note: _FixEmptyStructArrays is intentionally kept — Daft internally cannot handle empty StructArrays, not just an arrow2 FFI issue. ## Related Issues Closes #6347 | 6 个月前 | |
chore(tests): migrate internal usages of daft.udf to cls/func (#6348) update tests to remove a bunch of deprecation warnings in CI/tests ## Changes Made <!-- Describe what changes were made and why. Include implementation details if necessary. --> ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> | 6 个月前 | |
fix(native): preserve df metrics after stats manager shutdown (#6555) ## Changes Made The metrics were being collected correctly, but they were lost during teardown. take_input_snapshot(input_id) was called after RuntimeStatsManager had already exited, so the request failed and the executor returned an empty ExecutionStats. This change preserves finalized per-input metrics before the stats manager shuts down and makes them available through RuntimeStatsManagerHandle afterward. take_input_snapshot(input_id) now returns that stored final snapshot during executor teardown, and removes it once consumed so completed inputs do not accumulate indefinitely. | 5 个月前 | |
feat(swordfish): plan caching (#6278) ## Summary Introduces fingerprint-based plan caching so multiple executions of the same logical plan share a single pipeline, rather than building a new one each time. ### Key changes - ** PipelineMessage enum** — Replaces raw MicroPartition in all pipeline channels with Morsel { input_id, partition } and Flush(input_id), allowing a single pipeline to multiplex data from multiple logical inputs - **NativeExecutor becomes stateful** — Maintains a HashMap<u64, PlanState> keyed by plan fingerprint; repeated calls to run() with the same fingerprint reuse the existing pipeline - **MessageRouter** — Routes pipeline output to per-input_id unbounded channels so each caller gets only its own results - **try_finish() lifecycle API** — Callers signal input completion and collect per-input stats; pipeline is torn down when the last input finishes - **Per-input runtime stats** — RuntimeStatsManager tracks (NodeID, InputId) pairs with merge support for aggregated views - **next_event helper** — Shared tokio::select! loop abstraction used by blocking sinks, streaming sinks, intermediate ops, and joins - **Simplified BlockingSink trait** — Removed BlockingSinkFinalizeOutput::HasMoreOutput (was dead code); finalize now returns Vec<MicroPartition> directly; make_state takes InputId - **Simplified StreamingSink trait** — Removed StreamingSinkOutput::HasMoreOutput variant - **All node types updated** — Sources, intermediate ops, blocking sinks, streaming sinks, joins, and concat all handle Morsel/Flush per-input routing - **Python integration** — native_executor.py and flotilla.py call try_finish() for lifecycle management ### Files changed ~49 files across daft-local-execution (pipeline, sources, sinks, joins, stats), daft-distributed, common-metrics, and Python runners. ## Test plan - [ ] Full test suite with DAFT_RUNNER=native - [ ] Full test suite with DAFT_RUNNER=ray - [ ] Concurrent query execution sharing the same plan fingerprint - [ ] Verify per-input stats are correctly reported and merged - [ ] TPC-H benchmarks to check for performance regressions 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com> Co-authored-by: R. C. Howell <5731503+rchowell@users.noreply.github.com> | 5 个月前 | |
fix: Preserve materialized stats when splitting granular projections (#6376) ## Changes Made <!-- Describe what changes were made and why. Include implementation details if necessary. --> Fix a painc error occurs where executing explain on a dataframe containing join and download operators in Ray Runner. For example: python df1 = daft.from_pydict( { "name": ["a", "b"], "url": ["https://www.daft.ai/"] * 2, } ).with_column("bytes", col("url").download()) df2 = daft.from_pydict( { "name": ["a", "b"], "value": [1, 2], } ) df = df1.join(df2, on="name", prefix="df2_") df.explain(True) ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> Signed-off-by: plotor <zhenchao.wang@hotmail.com> | 6 个月前 | |
feat(flotilla): Hash Join (#4700) | 1 年前 | |
feat(flotilla): Distributed Pivot (#5199) | 1 年前 | |
feat: No truncate in .collect preview (#5632) ## Changes Made Follow on from https://github.com/Eventual-Inc/Daft/pull/5597. Remove truncation logic when previewing with .collect(num_preview_rows=N) ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> ## Checklist - [ ] Documented in API Docs (if applicable) - [ ] Documented in User Guide (if applicable) - [ ] If adding a new documentation page, doc is added to docs/mkdocs.yml navigation - [ ] Documentation builds and is formatted properly | 10 个月前 | |
feat: creates daft-ir and daft-proto (#4539) ## Changes Made * Creates a daft-ir crate for a consolidate IR view + from/to proto implementations (WIP). * Creates a daft-proto crate which holds daft's IR protobuf definitions. * Adds initial schema, datatype, expression (rex), and logical plan (rel) definitions. ## TODOs. * The exports (and structure) are not finalized. * Effectively all relational operator (rel) definitions. * Testing. testing. testing. ## Related Issues * n/a ## Checklist - [n/a] Documented in API Docs (if applicable) - [n/a] Documented in User Guide (if applicable) - [n/a] If adding a new documentation page, doc is added to docs/mkdocs.yml navigation - [n/a] Documentation builds and is formatted properly (tag @/ccmao1130 for docs review) | 1 年前 | |
feat: Add a random shuffle operation (#6481) ## Changes Made Implement a global row-rearraging shuffle operation, similar to https://huggingface.co/docs/datasets/process#shuffle. All it does is generate a random integer per row and then sort by that integer. Rather than implementing a custom distributed and local operation, it is all done within expressions and the optimizer. Note, I still made it an explicit logical node because: 1. Its easier to understand in explain outputs 2. I don't think we should always assume that it optimizes the same way that sort does. For example, it might be reasonable to push limits through a shuffle, but not allowed for sorts. ## Related Issues First half of #2612. Second half to add a batch_size version. --------- Co-authored-by: greptile-apps[bot] <165735046+greptile-apps[bot]@users.noreply.github.com> | 5 个月前 | |
feat: implement _repr_mimebundle_ for DataFrame (#6385) # Problem DataFrame does not implement _repr_mimebundle_, so some notebook frontends cannot reliably render rich output. # Root Cause Only __repr__ and _repr_html_ are implemented today; the IPython mimebundle display protocol is missing. # Solution - Added _repr_mimebundle_(include, exclude) to DataFrame. - Returns both text/plain and text/html by default. - Supports include/exclude filtering. # Tests - Added two test groups in tests/dataframe/test_repr.py. - Ran -k mimebundle: 8 passed. # Impact Improves DataFrame display compatibility across notebook frontends without affecting query behavior. --------- Co-authored-by: Lord of Abyss <pancx@chinatelecom.cn> | 6 个月前 | |
feat: Streaming sample by size (#5663) ## Changes Made Implement a streaming sample by size algorithm for swordfish. - For sampling without replacement, we assign a random number to each row, and simply keep around the rows with the lowest N random numbers. - For sampling with replacement where size=N, we create N slots. For each slot, assign a random number to each row. If the random number is smaller than the existing one in that slot, replace it. This PR only implements for swordfish. For flotilla, we will need to do some merging of local samples, this will be in the next PR. Essentially, we keep the random numbers per row / slot, merge the data onto a single node, and then pick the N rows with the smallest random numbers (per slot if with replacement). ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> | 10 个月前 | |
feat: supports passing a projection kwargs in select (#4884) ## Changes Made Adds support for a dict or projections to select which is backwards compatible but consistent with .with_columns. This format makes it easy to list down the columns prior to inserting somewhere else. python In [1]: import daft In [2]: df = daft.from_pydict({"x": [1, 2, 3], "y": [4, 5, 6], "z": [7, 8, 9]}) In [3]: df = df.select({ ...: "sum": df["x"] + df["y"] + df["z"], ...: "product": df["x"] * df["y"] * df["z"], ...: }) In [4]: df.show() ╭───────┬─────────╮ │ sum ┆ product │ │ --- ┆ --- │ │ Int64 ┆ Int64 │ ╞═══════╪═════════╡ │ 12 ┆ 28 │ ├╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌┤ │ 15 ┆ 80 │ ├╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌┤ │ 18 ┆ 162 │ ╰───────┴─────────╯ (Showing first 3 of 3 rows) ## Related Issues Closes #4883 ## Checklist - [ ] Documented in API Docs (if applicable) - [ ] Documented in User Guide (if applicable) - [ ] If adding a new documentation page, doc is added to docs/mkdocs.yml navigation - [ ] Documentation builds and is formatted properly (tag @/ccmao1130 for docs review) | 1 年前 | |
feat: Add global agg support for Dataframe select (#6067) ## Changes Made We encountered an issue about aggregate functions are not permitted in the SELECT statement unless a GROUP BY clause is specified. for example: from daft import col import daft df = daft.from_pydict({"a": [1, 2, 3], "b": [4, 5, 6]}) df = df.select( col('a').sum().alias('sum') ) df.show() DaftCoreException: DaftError::ValueError Aggregation expressions are currently only allowed in agg, pivot, and window: sum(col(a)) This issue has been resolved in daft sql https://github.com/Eventual-Inc/Daft/pull/2799, but it is still not supported in dataframe. This PR is precisely aimed at addressing this issue. ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> I found that a related issue already exists. https://github.com/Eventual-Inc/Daft/issues/1979 --------- Co-authored-by: cancai <caican@xiaomi.com> | 6 个月前 | |
fix(cast): handle whitespace in string-to-number casting (#5955) ## Changes Made - Handle whitespace trimming for Utf8 -> numeric casting in cast.rs - Add Rust unit tests for whitespace handling in cast operations - Simplify run_process to return raw subprocess stdout - Remove | tr -d ' ' workaround from tests since cast now handles whitespace ## Motivation String-to-number casting in Daft was strict about whitespace, returning None for inputs like " 3 ". This differed from: - Python's int() - Pandas - NumPy - DuckDB All of which handle leading/trailing whitespace gracefully. ### Comparison (Numeric) | Engine | ' 3 ' → int | |--------|-----------------| | Python int() | 3 | | Pandas | 3 | | DuckDB | 3 | | Polars | None | | Daft (before) | None | | **Daft (after)** | **3** | ## Implementation The fix is implemented in src/daft-core/src/array/ops/cast.rs: - Before calling arrow2's cast, check if we're casting Utf8 -> numeric - If so, create a trimmed copy of the string array - Pass the trimmed data to arrow2's existing cast logic Trade-off: Temporarily duplicates string data during cast, but gains code reuse and maintainability. ## Scope This PR only affects **numeric** casting (int, float). ### Date Parsing Date parsing already handles whitespace correctly - chrono::NaiveDate::from_str (used by arrow2) is lenient with whitespace. Added a test to document this existing behavior. ## Related Issues Fixes the root cause of #5942 Related to #5738 | 8 个月前 | |
feat: union ops (#3872) adds the following to both dataframe and sql - union (distinct union) - union_all - union_by_name (distinct union by name) - union_all_by_name closes https://github.com/Eventual-Inc/Daft/issues/3855 | 1 年前 | |
fix: deflake a probabilistic test that is causing release uploads to fail (#6475) ## Changes Made This increases the number of files we generate to check sharding fairness which should reliably increase our fairness check. With num_files=100 this was flaky since some randomness could exceed our threshold. ## Related Issues https://github.com/Eventual-Inc/Daft/actions/runs/23474792750 | 6 个月前 | |
fix: preserve show options without setting format (#6102) ## Changes Made - Pass in format=None to Rust preview binding and use default PreviewFormat in that case. - Make table_display.rs single source of truth for configuring and printing tables. - Check PYTEST_CURRENT_TEST instead of custom DAFT_BOLD_TABLE_HEADERSflag during tests. ## Related Issues Closes #4147. | 6 个月前 | |
chore: Move Flight Server to Rust (#6519) ## Changes Made Provide in Rust so we have direct access to it for optimizations | 5 个月前 | |
fix: canonicalize negative NaN in multi-column sort comparator (#6215) ## Changes Made Fixes one of the bugs causing the nightly ["Run property based tests with Hypothesis"](https://github.com/Eventual-Inc/Daft/actions/workflows/property-based-tests.yml) workflow to fail since Feb 2: negative NaN sorts before regular values in multi-column sort. Failing runs caused by this bug: - [Feb 10 (ray job)](https://github.com/Eventual-Inc/Daft/actions/runs/21880931284) - [Feb 16 (both jobs)](https://github.com/Eventual-Inc/Daft/actions/runs/22076216982) ### Root cause build_multi_array_bicompare uses arrow-rs make_comparator for secondary sort columns. arrow-rs follows IEEE 754 total ordering where **negative NaN sorts before all values** (-NaN < -Inf < ... < +Inf < +NaN). Daft's sort contract treats **all NaN as greater than regular values** regardless of sign. Single-column sort already handles this correctly via cmp_float. The bug only affects multi-column sort when a float column is a secondary key and contains a negative NaN (which hypothesis floats() can generate). ### Fix Canonicalize negative NaN to positive NaN in float arrays before passing them to make_comparator. This is a small allocation trade-off (one new array per float column per sort) in exchange for a simple, targeted fix. A zero-allocation alternative would be to replace make_comparator with a custom comparator using cmp_float for float types, but that adds more code complexity. ### Verification - Unit test: test_negative_nan_multicolumn_sort fails before fix, passes after (with a contrasting test_negative_nan_single_column_sort that already passes) - Hypothesis: HYPOTHESIS_MAX_EXAMPLES=10000 --hypothesis-seed=0 passes locally (same config as nightly CI) | 7 个月前 | |
| 6 个月前 | ||
feat: string casing functions (#6096) ## Changes Made Use [heck](https://crates.io/crates/heck) crate to convert cols and text between string casings. ## Related Issues Closes #2550. | 7 个月前 | |
refactor: Cleanup Dtype Names (#5400) | 11 个月前 | |
feat: add current_date, current_timestamp, current_timezone and SQL aliases (#6495) ## Summary - Implements current_date, current_timestamp, and current_timezone from #3798 - Adds Rust scalar UDFs (CurrentDate, CurrentTimestamp, CurrentTimezone) in daft-functions-temporal - Registers SQL functions in the SQL planner - Adds Python API functions - 4 tests covering type correctness, value correctness, and SQL integration Closes partially #3798 ## Test plan - [x] All temporal tests pass (DAFT_RUNNER=native pytest tests/dataframe/test_temporals.py -v) - [x] ruff check passes on all changed Python files - [x] make build compiles successfully | 6 个月前 | |
fix: close connections in io tests and address other warnings (#6470) ## Changes Made * Iceberg tests weren't closing the sqlite connection which made our CI logs flooded with warnings. * Several other io tests had the same issue. * Also fixes some additional warnings which are making CI logs noisy. ## Related Issues N/A | 6 个月前 | |
[FEAT] [New Query Planner] All functional tests pass + add to CI. (#1274) This PR ensures that all functional tests pass for the new query planner and adds the new query planner to the CI job matrix; all tests are covered and pass except for the Python query planner optimization tests. The use_new_planner test fixture is removed, since DAFT_NEW_QUERY_PLANNER=1 make test will now work with expected test coverage. | 3 年前 | |
chore: Upgrade Ruff ruleset to 3.9 and add from __future__ import annotations (#4393) | 1 年前 | |
chore: Remove expression namespaces (#5619) ## Changes Made Bye bye ## Related Issues <!-- Link to related GitHub issues, e.g., "Closes #123" --> ## Checklist - [ ] Documented in API Docs (if applicable) - [ ] Documented in User Guide (if applicable) - [ ] If adding a new documentation page, doc is added to docs/mkdocs.yml navigation - [ ] Documentation builds and is formatted properly | 10 个月前 | |
[FEAT] Enable concat for swordfish (#2976) Implement concat as a streaming sink in swordfish. Concatenating two dataframes in a streaming executor is pretty simple. You just stream the left side, then stream the right side. This PR implements just that. --------- Co-authored-by: Colin Ho <colinho@Colins-MBP.localdomain> Co-authored-by: Colin Ho <colinho@Colins-MacBook-Pro.local> | 1 年前 | |
chore: Upgrade Ruff ruleset to 3.9 and add from __future__ import annotations (#4393) | 1 年前 | |
feat: add uuid function (#5983) ## Changes Made add uuid function python df = df.with_column("u1", uuid()) <!-- Describe what changes were made and why. Include implementation details if necessary. --> ## Related Issues Closes #3706 <!-- Link to related GitHub issues, e.g., "Closes #123" --> --------- Co-authored-by: wangzheyan <wangzheyan@bytedance.com> | 7 个月前 | |
| 6 个月前 | ||
chore: Upgrade Ruff ruleset to 3.9 and add from __future__ import annotations (#4393) | 1 年前 | |
[BUG] with_column with existing column name should not reorder columns (#2381) | 2 年前 | |
feat: with_column(s)_renamed expression for DataFrame (#3732) [issue #3705](https://github.com/Eventual-Inc/Daft/issues/3705) Add with_column_renamed and with_columns_renamed to DataFrame API (no SQL). def with_column_renamed(self, existing: str, new: str) -> "DataFrame": """Renames a column in the current DataFrame. Example: >>> import daft >>> df = daft.from_pydict({"x": [1, 2, 3], "y": [4, 5, 6]}) >>> df.with_column_renamed("x", "foo").show() ╭───────┬───────╮ │ foo ┆ y │ │ --- ┆ --- │ │ Int64 ┆ Int64 │ ╞═══════╪═══════╡ │ 1 ┆ 4 │ ├╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌┤ │ 2 ┆ 5 │ ├╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌┤ │ 3 ┆ 6 │ ╰───────┴───────╯ <BLANKLINE> (Showing first 3 of 3 rows) Args: existing (str): name of the existing column to rename new (str): new name for the column Returns: DataFrame: DataFrame with the column renamed. """ def with_columns_renamed(self, colsMap: Dict[str, str]) -> "DataFrame": """Renames multiple columns in the current DataFrame. Example: >>> import daft >>> df = daft.from_pydict({"x": [1, 2, 3], "y": [4, 5, 6]}) >>> df.with_columns_renamed({"x": "foo", "y": "bar"}).show() ╭───────┬───────╮ │ foo ┆ bar │ │ --- ┆ --- │ │ Int64 ┆ Int64 │ ╞═══════╪═══════╡ │ 1 ┆ 4 │ ├╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌┤ │ 2 ┆ 5 │ ├╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌┤ │ 3 ┆ 6 │ ╰───────┴───────╯ <BLANKLINE> (Showing first 3 of 3 rows) """ | 1 年前 | |
[BUG] with_column with existing column name should not reorder columns (#2381) | 2 年前 | |
feat: with_column(s)_renamed expression for DataFrame (#3732) [issue #3705](https://github.com/Eventual-Inc/Daft/issues/3705) Add with_column_renamed and with_columns_renamed to DataFrame API (no SQL). def with_column_renamed(self, existing: str, new: str) -> "DataFrame": """Renames a column in the current DataFrame. Example: >>> import daft >>> df = daft.from_pydict({"x": [1, 2, 3], "y": [4, 5, 6]}) >>> df.with_column_renamed("x", "foo").show() ╭───────┬───────╮ │ foo ┆ y │ │ --- ┆ --- │ │ Int64 ┆ Int64 │ ╞═══════╪═══════╡ │ 1 ┆ 4 │ ├╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌┤ │ 2 ┆ 5 │ ├╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌┤ │ 3 ┆ 6 │ ╰───────┴───────╯ <BLANKLINE> (Showing first 3 of 3 rows) Args: existing (str): name of the existing column to rename new (str): new name for the column Returns: DataFrame: DataFrame with the column renamed. """ def with_columns_renamed(self, colsMap: Dict[str, str]) -> "DataFrame": """Renames multiple columns in the current DataFrame. Example: >>> import daft >>> df = daft.from_pydict({"x": [1, 2, 3], "y": [4, 5, 6]}) >>> df.with_columns_renamed({"x": "foo", "y": "bar"}).show() ╭───────┬───────╮ │ foo ┆ bar │ │ --- ┆ --- │ │ Int64 ┆ Int64 │ ╞═══════╪═══════╡ │ 1 ┆ 4 │ ├╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌┤ │ 2 ┆ 5 │ ├╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌┤ │ 3 ┆ 6 │ ╰───────┴───────╯ <BLANKLINE> (Showing first 3 of 3 rows) """ | 1 年前 |
| 文件 | 最后提交记录 | 最后更新时间 |
|---|---|---|
| 3 年前 | ||
| 3 年前 | ||
| 9 个月前 | ||
| 6 个月前 | ||
| 7 个月前 | ||
| 1 年前 | ||
| 10 个月前 | ||
| 1 年前 | ||
| 1 年前 | ||
| 1 年前 | ||
| 6 个月前 | ||
| 6 个月前 | ||
| 1 年前 | ||
| 7 个月前 | ||
| 5 个月前 | ||
| 7 个月前 | ||
| 7 个月前 | ||
| 1 年前 | ||
| 1 年前 | ||
| 1 年前 | ||
| 5 个月前 | ||
| 1 年前 | ||
| 6 个月前 | ||
| 7 个月前 | ||
| 8 个月前 | ||
| 6 个月前 | ||
| 6 个月前 | ||
| 5 个月前 | ||
| 5 个月前 | ||
| 6 个月前 | ||
| 1 年前 | ||
| 1 年前 | ||
| 10 个月前 | ||
| 1 年前 | ||
| 5 个月前 | ||
| 6 个月前 | ||
| 10 个月前 | ||
| 1 年前 | ||
| 6 个月前 | ||
| 8 个月前 | ||
| 1 年前 | ||
| 6 个月前 | ||
| 6 个月前 | ||
| 5 个月前 | ||
| 7 个月前 | ||
| 6 个月前 | ||
| 7 个月前 | ||
| 11 个月前 | ||
| 6 个月前 | ||
| 6 个月前 | ||
| 3 年前 | ||
| 1 年前 | ||
| 10 个月前 | ||
| 1 年前 | ||
| 1 年前 | ||
| 7 个月前 | ||
| 6 个月前 | ||
| 1 年前 | ||
| 2 年前 | ||
| 1 年前 | ||
| 2 年前 | ||
| 1 年前 |