Skip to content
Merged
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
18 changes: 15 additions & 3 deletions src/altertable_flightsql/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -697,7 +697,18 @@ def _get_parameter_as_pyarrow(
elif isinstance(parameters, pa.RecordBatch):
return parameters
elif isinstance(parameters, Mapping):
return pa.record_batch({key: [value] for (key, value) in parameters.items()})
return pa.RecordBatch.from_pydict(
{
key: (
pa.array([value], type=self._parameter_schema.field(key).type)
if value is None
and self._parameter_schema is not None
and key in self._parameter_schema.names
else [value]
)
for key, value in parameters.items()
}
)
elif isinstance(parameters, Sequence):
if self._parameter_schema is None:
raise ValueError(
Expand All @@ -711,10 +722,11 @@ def _get_parameter_as_pyarrow(
f"Expected {len(self._parameter_schema)} parameters, but got {len(parameters)}"
)
param_dict = {
field.name: [value] for field, value in zip(self._parameter_schema, parameters)
field.name: pa.array([value], type=field.type) if value is None else [value]
for field, value in zip(self._parameter_schema, parameters)
}

return pa.record_batch(param_dict)
return pa.RecordBatch.from_pydict(param_dict)
else:
raise TypeError(
f"Unsupported parameter type: {type(parameters)}. "
Expand Down
21 changes: 20 additions & 1 deletion tests/test_client.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
from types import SimpleNamespace

import pyarrow as pa
import pyarrow.flight as flight
import pytest
from google.protobuf import any_pb2

from altertable_flightsql.client import Client
from altertable_flightsql.client import Client, PreparedStatement
from altertable_flightsql.generated import arrow_flight_pb2 as flight_pb2


Expand Down Expand Up @@ -69,6 +70,24 @@ def _action_body_bytes(action) -> bytes:
return bytes(body)


@pytest.mark.parametrize(
"values",
[
[None],
{"amount": None},
],
ids=["positional", "mapping"],
)
def test_python_null_parameters_use_prepared_type(values):
parameter_schema = pa.schema([("amount", pa.float64())])
statement = PreparedStatement(None, b"handle", parameter_schema=parameter_schema)

parameters = statement._get_parameter_as_pyarrow(values)

assert parameters.schema.equals(parameter_schema)
assert parameters.to_pydict() == {"amount": [None]}


def test_set_options_serializes_flight_session_request_without_any():
flight_client = FakeFlightClient()
client = _client_backed_by(flight_client)
Expand Down
Loading