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
3 changes: 2 additions & 1 deletion .github/workflows/python-app.yml
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ jobs:
- name: Install dependencies
run: |
python -m pip install --upgrade pip
pip install flake8 pytest
pip install flake8 pytest ipython
if [ -f requirements.txt ]; then pip install -r requirements.txt; fi
- name: Lint with flake8
run: |
Expand All @@ -40,5 +40,6 @@ jobs:
run: |
# run tests independently to enable tests using spark session
for f in tests/test_*.py; do
echo "Running tests in $f"
pytest -q "$f" || exit 1
done
26 changes: 19 additions & 7 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,13 @@ A RumbleSession is a wrapper around a SparkSession that additionally makes sure

JSONiq queries are invoked with rumble.jsoniq() in a way similar to the way Spark SQL queries are invoked with spark.sql().

Use `rumble.xquery()` to run queries with XQuery 3.1 as the default language. It accepts the same keyword variable bindings and returns the same `SequenceOfItems` as `rumble.jsoniq()`, while preserving the session configuration. An explicit language version declaration in the query overrides the default.

```python
result = rumble.xquery("<greeting>{$name}</greeting>", name="World")
print(result.serialize()) # XML declaration followed by <greeting>World</greeting>
```

JSONiq variables can be bound to lists of JSON values (str, int, float, True, False, None, dict, list) or to Pyspark DataFrames. A JSONiq query can use as many variables as needed (for example, it can join between different collections).

It will later also be possible to read tables registered in the Hive metastore, similar to spark.sql(). Alternatively, the JSONiq query can also read many files of many different formats from many places (local drive, HTTP, S3, HDFS, ...) directly with simple builtin function calls such as json-lines(), text-file(), parquet-file(), csv-file(), etc. See [RumbleDB's documentation](https://docs.rumbledb.org/writing-jsoniq-queries-in-python).
Expand All @@ -40,6 +47,13 @@ It is also possible to write the sequence of items to the local disk, to HDFS, t

The library also contains a jsoniq magic that allows you to directly write JSONiq queries in a Jupyter notebook and see the results automatically output on the screen. In notebooks, you can use `%%jsoniq -s` (or `--serialize`) to serialize the result sequence to text output, according to the XSLT and XQuery Serialization 3.1 specification by W3C. The method and serialization options can all be specified in the query with option declarations, following the XQuery/JSONiq standard. `%%jsoniq -j` shows the results in JSON lines format, while `%%jsoniq -pdf` shows a pandas data frame, and `%%jsoniq -df` shows a Spark data frame.

The same notebook extension also registers `%%xquery`, which calls `rumble.xquery()` and defaults to serialized output (`-s`) with the XML serialization method. Query option declarations in the `http://www.w3.org/2010/xslt-xquery-serialization` namespace can override the serialization method. It accepts the same options: use `-j`, `-df`, or `-pdf` to select another output format, `-u` to apply updates, and `-t` to show execution time.

```xquery
%%xquery
<greeting>Hello, world!</greeting>
```

The design goal is that it is possible to chain DataFrames between JSONiq and Spark SQL queries seamlessly. For example, JSONiq can be used to clean up very messy data and turn it into a clean DataFrame, which can then be processed with Spark SQL, spark.ml, etc.

Any feedback or error reports are very welcome.
Expand Down Expand Up @@ -374,18 +388,16 @@ Even more queries can be found [here](https://colab.research.google.com/github/R

## Version 3.0.0
- Upgraded to RumbleDB 3.0.0 and its immutable configuration and external bindings APIs.
- Configuration reads use the Java API directly, such as `getInt("runtime.resultsSizeCap")` and `getBoolean("debug.showErrorInfo")`. The Python `set(path, value)` helper rebuilds the immutable Java configuration. Configuration changes apply to subsequent queries; existing sequences retain their compilation settings.
- Fixed object conversion and binding query results as DataFrames. Keyword bindings are scoped to a query and restore persistent bindings even when query compilation fails.

Configuration can also be changed using RumbleDB 3.0's dot-separated paths:

- There is a breaking change in how configuration parameters are set. New:
```python
rumble.getRumbleConf().set("runtime.resultsSizeCap", 100)
rumble.getRumbleConf().set("runtime.materializationCap", 100000)
rumble.getRumbleConf().set("debug.showErrorInfo", True)
```

The result size cap controls `first()` and notebook display. `json()` retrieves all items, subject to the separate materialization cap.
- Calling old configuration functions leads to a message explaining the new syntax. Configuration changes apply to subsequent queries; existing sequences retain their compilation settings.
- Fixed object conversion and binding query results as DataFrames. Keyword bindings are scoped to a query and restore persistent bindings even when query compilation fails.
- New rumble.xquery() call and %%xquery magic.
- In the %%jsoniq magic, the new -s parameter outputs the results as a single string following the W3C Serialization 3.1 specification. If using `rumble.xquery()` or `rumble.jsoniq()` call, the same can be achieved by chaining a `.serialize() ` call returning a string. All standard serialization methods (xml, json, xhtml, html, text, adaptive) are available. For %%xquery, this is the default behavior.

## Version 2.1.9
- Fixed a bug in the inferred conversion to DataFrames of output involving arrays of objects.
Expand Down
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"

[project]
name = "jsoniq"
version = "3.0.0"
version = "3.0.1"
description = "Python edition of RumbleDB, a JSONiq engine"
requires-python = ">=3.11"
dependencies = [
Expand Down
12 changes: 11 additions & 1 deletion src/jsoniq/session.py
Original file line number Diff line number Diff line change
Expand Up @@ -287,14 +287,24 @@ def bindDataFrameAsVariable(self, name: str, df):
return self;

def jsoniq(self, str, **kwargs):
return self._run_query(str, self._jrumblesession, kwargs)

def xquery(self, str, **kwargs):
"""Run a query with XQuery 3.1 as the default language."""
builder = self._jrumblesession.getConfiguration().toBuilder()
configuration = getattr(builder, "with")("semantics.queryLanguage", "xquery31").build()
engine = self._sparksession._jvm.org.rumbledb.api.Rumble(configuration)
return self._run_query(str, engine, kwargs)

def _run_query(self, query, engine, kwargs):
previous_bindings = self._bindings.copy()
try:
for key, value in kwargs.items():
self.bind(f"${key}", value)
bindings = self._sparksession._jvm.org.rumbledb.api.ExternalBindings()
for name, (method, value) in self._bindings.items():
getattr(bindings, method)(name, value)
sequence = self._jrumblesession.runQuery(str, bindings)
sequence = engine.runQuery(query, bindings)
return SequenceOfItems(sequence, self)
finally:
self._bindings = previous_bindings
Expand Down
12 changes: 10 additions & 2 deletions src/jsoniqmagic/magic.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,17 +27,21 @@ class JSONiqMagic(Magics):
@argument(
'-u', '--apply-updates', action='store_true', help='Applies updates if a PUL is output.'
)
def run(self, line, cell=None, timed=False):
def run(self, line, cell=None, timed=False, language="jsoniq"):
if cell is None:
data = line
else:
data = cell

args = parse_argstring(self.run, line)
if language == "xquery" and not (
args.json or args.pandas_data_frame or args.pyspark_data_frame or args.apply_updates
):
args.serialize = True
start = time.time()
try:
rumble = RumbleSession.builder.getOrCreate();
response = rumble.jsoniq(data);
response = getattr(rumble, language)(data);
except Py4JJavaError as e:
print(e.java_exception.getMessage())
return
Expand Down Expand Up @@ -195,3 +199,7 @@ def run(self, line, cell=None, timed=False):
@cell_magic
def jsoniq(self, line, cell=None):
return self.run(line, cell, False)

@cell_magic
def xquery(self, line, cell=None):
return self.run(line, cell, False, language="xquery")
90 changes: 90 additions & 0 deletions tests/fixtures/legacy_configuration_methods.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,90 @@
{
"source": "https://rumbledb.org/docs/latest/api/org/rumbledb/config/RumbleRuntimeConfiguration.html",
"signatures": [
"getDefaultConfiguration()",
"getPort()",
"getHost()",
"getAllowedURIPrefixes()",
"setAllowedURIPrefixes(java.util.List)",
"getInputFormat()",
"setInputFormat(java.lang.String)",
"getNumberOfOutputPartitions()",
"setNumberOfOutputPartitions(int)",
"isCheckReturnTypeOfBuiltinFunctions()",
"setCheckReturnTypeOfBuiltinFunctions(boolean)",
"init()",
"getOverwrite()",
"getShowErrorInfo()",
"setShowErrorInfo(boolean)",
"getLaxJSONNullValidation()",
"setLaxJSONNullValidation(boolean)",
"getLogPath()",
"getQueryPath()",
"getOutputPath()",
"getQuery()",
"getShellFilter()",
"getNativeSQLPredicates()",
"getDataFrameExecutionModeDetection()",
"setLogPath(java.lang.String)",
"setQueryPath(java.lang.String)",
"setOutputPath(java.lang.String)",
"setShellFilter(java.lang.String)",
"setNativeSQLPredicates(boolean)",
"setDataFrameExecutionModeDetection(boolean)",
"getResultSizeCap()",
"setResultSizeCap(int)",
"getMaterializationCap()",
"setMaterializationCap(int)",
"getExternalVariablesReadFromDataFrames()",
"getExternalVariablesReadFromListsOfItems()",
"getExternalVariableValue(org.rumbledb.context.Name)",
"getUnparsedExternalVariableValue(org.rumbledb.context.Name)",
"getExternalVariableValueReadFromFile(org.rumbledb.context.Name)",
"getExternalVariableValueReadFromDataFrame(org.rumbledb.context.Name)",
"setExternalVariableValue(org.rumbledb.context.Name,org.apache.spark.sql.Dataset)",
"setExternalVariableValue(java.lang.String,org.apache.spark.sql.Dataset)",
"setExternalVariableValue(org.rumbledb.context.Name,java.util.List)",
"setExternalVariableValue(java.lang.String,java.util.List)",
"resetExternalVariableValue(org.rumbledb.context.Name)",
"resetExternalVariableValue(java.lang.String)",
"readFromStandardInput(org.rumbledb.context.Name)",
"getInputFormat(org.rumbledb.context.Name)",
"isShell()",
"isServer()",
"setPrintIteratorTree(boolean)",
"isPrintIteratorTree()",
"doStaticAnalysis()",
"printInferredTypes()",
"dateWithTimezone()",
"setDateWithTimezone(boolean)",
"parallelExecution()",
"setParallelExecution(boolean)",
"dataFrameExecution()",
"setDataFrameExecution(boolean)",
"nativeExecution()",
"setNativeExecution(boolean)",
"functionInlining()",
"setFunctionInlining(boolean)",
"applyUpdates()",
"setApplyUpdates(boolean)",
"optimizeGeneralComparisonToValueComparison()",
"setOptimizeGeneralComparisonToValueComparison(boolean)",
"getQueryLanguage()",
"setQueryLanguage(java.lang.String)",
"getStaticBaseUri()",
"setStaticBaseUri(java.lang.String)",
"optimizeSteps()",
"setOptimizeSteps(boolean)",
"optimizeStepExperimental()",
"setOptimizeStepsExperimental(boolean)",
"optimizeParentPointers()",
"setOptimizeParentPointers(boolean)",
"isLocal()",
"toString()",
"write(com.esotericsoftware.kryo.Kryo,com.esotericsoftware.kryo.io.Output)",
"read(com.esotericsoftware.kryo.Kryo,com.esotericsoftware.kryo.io.Input)",
"getXmlVersion()",
"setXmlVersion(java.lang.String)",
"getSerializationParameters()"
]
}
67 changes: 67 additions & 0 deletions tests/test_configuration_migration.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
"""Check all documented legacy signatures without starting Spark."""

import json
from pathlib import Path
import runpy

import pytest


ROOT = Path(__file__).resolve().parents[1]
RumbleConfiguration = runpy.run_path(
str(ROOT / "src/jsoniq/configuration.py")
)["RumbleConfiguration"]
SIGNATURES = json.loads(
(ROOT / "tests/fixtures/legacy_configuration_methods.json").read_text()
)["signatures"]


@pytest.mark.parametrize("signature", SIGNATURES)
def test_every_documented_signature_raises_migration_error(signature):
name, parameters = signature[:-1].split("(")
args = [object() for _ in parameters.split(",")] if parameters else []
# Any attempt to forward a legacy method to Java would fail on this object.
configuration = RumbleConfiguration(object())
assert name in RumbleConfiguration.__dict__
with pytest.raises(NotImplementedError) as error:
getattr(configuration, name)(*args)
assert f"{name}() was removed in RumbleDB 3.0.0." in str(error.value)
assert len(str(error.value).split("3.0.0. ", 1)[1]) > 0


@pytest.mark.parametrize("name,args", [
("getHost", ()), ("getPort", ()), ("isServer", ()),
("getAllowedURIPrefixes", ()), ("setAllowedURIPrefixes", (["file:/"],)),
])
def test_server_and_uri_messages_direct_users_to_notebooks(name, args):
with pytest.raises(NotImplementedError) as error:
getattr(RumbleConfiguration(object()), name)(*args)
message = str(error.value)
assert "removed" in message
assert "server feature" in message
assert "Python library" in message
assert "notebooks" in message


@pytest.mark.parametrize("name,args", [
("getInputFormat", ()), ("getInputFormat", ("input",)),
("setInputFormat", ("json",)), ("readFromStandardInput", ("input",)),
])
def test_input_messages_preserve_json_arrays_and_explain_api_scope(name, args):
with pytest.raises(NotImplementedError) as error:
getattr(RumbleConfiguration(object()), name)(*args)
message = str(error.value)
assert 'rumble.bindOne("$input", json.load(sys.stdin))' in message
assert 'rumble.bindOne("$input", sys.stdin.read())' in message
assert "does not expose them yet" in message


@pytest.mark.parametrize("name,value,command", [
("setResultSizeCap", 42, 'set("runtime.resultsSizeCap", 42)'),
("setShowErrorInfo", True, 'set("debug.showErrorInfo", True)'),
("setQueryPath", 'a"b.jq', 'set("input.queryPath", \'a"b.jq\')'),
])
def test_setter_messages_include_python_value(name, value, command):
with pytest.raises(NotImplementedError) as error:
getattr(RumbleConfiguration(object()), name)(value)
assert command in str(error.value)
95 changes: 95 additions & 0 deletions tests/test_magic_serialization.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
from unittest.mock import Mock, patch

import pytest

pytest.importorskip("IPython")

from jsoniq import RumbleSession
from jsoniqmagic import JSONiqMagic


@pytest.mark.parametrize("option", ["-s", "--serialize", "-s -j", "-s -df", "-s -pdf", "-s -u"])
@pytest.mark.parametrize("serialized", ['<result>hello "world"</result>\nsecond line', ""])
def test_magic_prints_serialized_sequence_directly(option, serialized, capsys):
session = Mock()
response = session.jsoniq.return_value
response.serialize.return_value = serialized

with patch.object(RumbleSession.Builder, "getOrCreate", return_value=session):
JSONiqMagic().jsoniq(option, "1 to 3")

session.jsoniq.assert_called_once_with("1 to 3")
response.serialize.assert_called_once_with()
response.take.assert_not_called()
response.df.assert_not_called()
response.pdf.assert_not_called()
response.applyPUL.assert_not_called()
assert capsys.readouterr().out == serialized + "\n"


def test_magic_reports_serialization_errors(capsys):
session = Mock()
session.jsoniq.return_value.serialize.side_effect = RuntimeError("Serialization failed")

with patch.object(RumbleSession.Builder, "getOrCreate", return_value=session):
JSONiqMagic().jsoniq("-s", "1 to 3")

assert "Serialization failed" in capsys.readouterr().out


def test_magic_serialization_preserves_timing(capsys):
session = Mock()
session.jsoniq.return_value.serialize.return_value = "hello"

with patch.object(RumbleSession.Builder, "getOrCreate", return_value=session):
JSONiqMagic().jsoniq("-s -t", '"hello"')

output = capsys.readouterr().out
assert output.startswith("hello\nResponse time: ")
assert output.endswith(" ms\n")


@pytest.mark.parametrize("option", ["", "-s", "--serialize", "-s -j", "-t"])
def test_xquery_magic_defaults_to_serialization(option, capsys):
session = Mock()
response = session.xquery.return_value
response.serialize.return_value = '<greeting>Hello "world"</greeting>'

with patch.object(RumbleSession.Builder, "getOrCreate", return_value=session):
JSONiqMagic().xquery(option, "<greeting>Hello</greeting>")

session.xquery.assert_called_once_with("<greeting>Hello</greeting>")
session.jsoniq.assert_not_called()
response.serialize.assert_called_once_with()
response.take.assert_not_called()
output = capsys.readouterr().out
assert output.startswith('<greeting>Hello "world"</greeting>\n')
if option == "-t":
assert "Response time:" in output


@pytest.mark.parametrize("option", ["-j", "-df", "-pdf", "-u"])
def test_xquery_magic_explicit_options_override_default(option, capsys):
session = Mock()
response = session.xquery.return_value
session.getRumbleConf.return_value.getInt.return_value = 10
item = Mock()
item.serializeAsJSON.return_value = "42"
response.take.return_value = [item]
response.availableOutputs.return_value = ["PUL"]

with patch.object(RumbleSession.Builder, "getOrCreate", return_value=session):
result = JSONiqMagic().xquery(option, "42")

response.serialize.assert_not_called()
if option == "-j":
assert capsys.readouterr().out == "42\n"
elif option == "-df":
response.df.return_value.show.assert_called_once_with()
response.take.assert_not_called()
elif option == "-pdf":
assert result is response.pdf.return_value
response.take.assert_not_called()
else:
response.applyPUL.assert_called_once_with()
assert "Updates applied successfully." in capsys.readouterr().out
Loading
Loading