|
| 1 | +import datetime |
| 2 | +import numpy as np |
| 3 | +import pandas as pd |
| 4 | +import pyarrow as pa |
| 5 | +import pyarrow.compute as pc |
| 6 | +import pytest |
| 7 | + |
| 8 | +from arcticdb.version_store.processing import QueryBuilder, ExpressionNode |
| 9 | +from arcticdb.util.test import assert_frame_equal |
| 10 | + |
| 11 | + |
| 12 | +def df_with_all_column_types(num_rows=100): |
| 13 | + data = { |
| 14 | + "int_col": np.arange(num_rows, dtype=np.int64), |
| 15 | + "float_col": [np.nan if i%20==5 else i for i in range(num_rows)], |
| 16 | + "str_col": [f"str_{i}" for i in range(num_rows)], |
| 17 | + "bool_col": [i%2 == 0 for i in range(num_rows)], |
| 18 | + "datetime_col": pd.date_range(start=pd.Timestamp(2025, 1, 1), periods=num_rows) |
| 19 | + } |
| 20 | + index = pd.date_range(start=pd.Timestamp(2025, 1, 1), periods=num_rows) |
| 21 | + return pd.DataFrame(data=data, index=index) |
| 22 | + |
| 23 | + |
| 24 | +def compare_against_pyarrow(pyarrow_expr_str, expected_adb_expr, lib, function_map = None, expect_equal=True): |
| 25 | + adb_expr = ExpressionNode.from_pyarrow_expression_str(pyarrow_expr_str, function_map) |
| 26 | + assert str(adb_expr) == str(expected_adb_expr) |
| 27 | + pa_expr = eval(pyarrow_expr_str) |
| 28 | + |
| 29 | + # Setup |
| 30 | + sym = "sym" |
| 31 | + df = df_with_all_column_types() |
| 32 | + lib.write(sym, df) |
| 33 | + pa_table = pa.Table.from_pandas(df) |
| 34 | + |
| 35 | + # Apply filter to adb |
| 36 | + q = QueryBuilder() |
| 37 | + q = q[adb_expr] |
| 38 | + adb_result = lib.read(sym, query_builder=q).data |
| 39 | + |
| 40 | + # Apply filter to pyarrow |
| 41 | + pa_result = pa_table.filter(pa_expr).to_pandas() |
| 42 | + |
| 43 | + if expect_equal: |
| 44 | + assert_frame_equal(adb_result, pa_result) |
| 45 | + else: |
| 46 | + assert len(adb_result) != len(pa_result) |
| 47 | + |
| 48 | + |
| 49 | +def test_basic_filters(lmdb_version_store_v1): |
| 50 | + lib = lmdb_version_store_v1 |
| 51 | + |
| 52 | + # Filter by boolean column |
| 53 | + expr = f"pc.field('bool_col')" |
| 54 | + expected_expr = ExpressionNode.column_ref('bool_col') |
| 55 | + compare_against_pyarrow(expr, expected_expr, lib) |
| 56 | + |
| 57 | + # Filter by comparison |
| 58 | + for op in ["<", "<=", "==", ">=", ">"]: |
| 59 | + expr = f"pc.field('int_col') {op} 50" |
| 60 | + expected_expr = eval(f"ExpressionNode.column_ref('int_col') {op} 50") |
| 61 | + compare_against_pyarrow(expr, expected_expr, lib) |
| 62 | + |
| 63 | + # Filter with unary operators |
| 64 | + expr = "~pc.field('bool_col')" |
| 65 | + expected_expr = ~ExpressionNode.column_ref('bool_col') |
| 66 | + compare_against_pyarrow(expr, expected_expr, lib) |
| 67 | + |
| 68 | + # Filter with binary operators |
| 69 | + for op in ["+", "-", "*", "/"]: |
| 70 | + expr = f"pc.field('float_col') {op} 5.0 < 50.0" |
| 71 | + expected_expr = eval(f"ExpressionNode.column_ref('float_col') {op} 5.0 < 50.0") |
| 72 | + compare_against_pyarrow(expr, expected_expr, lib) |
| 73 | + |
| 74 | + for op in ["&", "|"]: |
| 75 | + expr = f"pc.field('bool_col') {op} (pc.field('int_col') < 50)" |
| 76 | + expected_expr = eval(f"ExpressionNode.column_ref('bool_col') {op} (ExpressionNode.column_ref('int_col') < 50)") |
| 77 | + compare_against_pyarrow(expr, expected_expr, lib) |
| 78 | + |
| 79 | + # Filter with expression method calls |
| 80 | + expr = "pc.field('str_col').isin(['str_0', 'str_10', 'str_20'])" |
| 81 | + expected_expr = ExpressionNode.column_ref('str_col').isin(['str_0', 'str_10', 'str_20']) |
| 82 | + compare_against_pyarrow(expr, expected_expr, lib) |
| 83 | + |
| 84 | + expr = "pc.field('float_col').is_nan()" |
| 85 | + expected_expr = ExpressionNode.column_ref('float_col').isnull() |
| 86 | + # We expect a different result between adb and pyarrow because of the different nan/null handling |
| 87 | + compare_against_pyarrow(expr, expected_expr, lib, expect_equal=False) |
| 88 | + |
| 89 | + expr = "pc.field('float_col').is_null()" |
| 90 | + expected_expr = ExpressionNode.column_ref('float_col').isnull() |
| 91 | + compare_against_pyarrow(expr, expected_expr, lib) |
| 92 | + |
| 93 | + expr = "pc.field('float_col').is_valid()" |
| 94 | + expected_expr = ExpressionNode.column_ref('float_col').notnull() |
| 95 | + compare_against_pyarrow(expr, expected_expr, lib) |
| 96 | + |
| 97 | +def test_complex_filters(lmdb_version_store_v1): |
| 98 | + lib = lmdb_version_store_v1 |
| 99 | + |
| 100 | + # Nested complex filters |
| 101 | + expr = "((pc.field('float_col') * 2) > 20.0) & (pc.field('int_col') <= pc.scalar(60)) | pc.field('bool_col')" |
| 102 | + expected_expr = (ExpressionNode.column_ref('float_col') * 2 > 20.0) & (ExpressionNode.column_ref('int_col') <= 60) | ExpressionNode.column_ref('bool_col') |
| 103 | + compare_against_pyarrow(expr, expected_expr, lib) |
| 104 | + |
| 105 | + expr = "((pc.field('float_col') / 2) > 20.0) & (pc.field('float_col') <= pc.scalar(60)) & pc.field('str_col').isin(['str_30', 'str_41', 'str_42', 'str_53', 'str_99'])" |
| 106 | + expected_expr = (ExpressionNode.column_ref('float_col') / 2 > 20.0) & (ExpressionNode.column_ref('float_col') <= 60) & ExpressionNode.column_ref('str_col').isin(['str_30', 'str_41', 'str_42', 'str_53', 'str_99']) |
| 107 | + compare_against_pyarrow(expr, expected_expr, lib) |
| 108 | + |
| 109 | + # Filters with function calls |
| 110 | + function_map = { |
| 111 | + "datetime.datetime": datetime.datetime, |
| 112 | + "abs": abs, |
| 113 | + } |
| 114 | + expr = "pc.field('datetime_col') < datetime.datetime(2025, 1, 20)" |
| 115 | + expected_expr = ExpressionNode.column_ref('datetime_col') < datetime.datetime(2025, 1, 20) |
| 116 | + compare_against_pyarrow(expr, expected_expr, lib, function_map) |
| 117 | + |
| 118 | + expr = "(pc.field('datetime_col') < datetime.datetime(2025, 1, abs(-20))) & (pc.field('int_col') >= abs(-5))" |
| 119 | + expected_expr = (ExpressionNode.column_ref('datetime_col') < datetime.datetime(2025, 1, abs(-20))) & (ExpressionNode.column_ref('int_col') >= abs(-5)) |
| 120 | + compare_against_pyarrow(expr, expected_expr, lib, function_map) |
| 121 | + |
| 122 | +def test_broken_filters(): |
| 123 | + # ill-formated filter |
| 124 | + expr = "pc.field('float_col'" |
| 125 | + with pytest.raises(ValueError): |
| 126 | + ExpressionNode.from_pyarrow_expression_str(expr) |
| 127 | + |
| 128 | + # pyarrow expressions only support single comparisons |
| 129 | + expr = "1 < pc.field('int_col') < 10" |
| 130 | + with pytest.raises(ValueError): |
| 131 | + ExpressionNode.from_pyarrow_expression_str(expr) |
| 132 | + |
| 133 | + # calling a mising function |
| 134 | + expr = "some.missing.function(5)" |
| 135 | + with pytest.raises(ValueError): |
| 136 | + ExpressionNode.from_pyarrow_expression_str(expr) |
0 commit comments