This is an automated email from the ASF dual-hosted git repository.
bossenti pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to refs/heads/dev by this push:
new f71bd9629 [#1254]: extend data lake measure endpoints get method to
process query parameters (#1349)
f71bd9629 is described below
commit f71bd96290b3f9068884a32966aa9f41de49b0e9
Author: Tim <[email protected]>
AuthorDate: Mon Feb 27 18:47:42 2023 +0100
[#1254]: extend data lake measure endpoints get method to process query
parameters (#1349)
* feature(#1254): add query parameters for get endpoint for measurements
Signed-off-by: bossenti <[email protected]>
* feature(#1254): add examples for query parameters
Signed-off-by: bossenti <[email protected]>
* chore: update example notebooks
Signed-off-by: bossenti <[email protected]>
* chore: add missing update of unit tests
Signed-off-by: bossenti <[email protected]>
* chore: add missing file headers
Signed-off-by: bossenti <[email protected]>
* chore: fix typo
Signed-off-by: bossenti <[email protected]>
* chore: fix typos
Signed-off-by: bossenti <[email protected]>
* feature(#1254): add examples to docs
Signed-off-by: bossenti <[email protected]>
* chore: add demo for initial setup
Signed-off-by: bossenti <[email protected]>
* chore: fix formatting
Signed-off-by: bossenti <[email protected]>
---------
Signed-off-by: bossenti <[email protected]>
---
streampipes-client-python/README.md | 2 +-
...introduction-to-streampipes-python-client.ipynb | 80 ++++++--
...cting-data-from-the-streampipes-data-lake.ipynb | 153 +++++++++++++-
...ive-data-from-the-streampipes-data-stream.ipynb | 12 ++
...ine-learning-on-a-streampipes-data-stream.ipynb | 13 +-
.../docs/img/tutorial-preparation.gif | Bin 0 -> 2569748 bytes
.../streampipes/client/client.py | 2 +-
.../streampipes/endpoint/api/data_lake_measure.py | 220 ++++++++++++++++++++-
.../streampipes/endpoint/api/data_stream.py | 2 +-
.../streampipes/endpoint/endpoint.py | 2 +-
.../tests/client/test_data_lake_series.py | 2 +-
.../tests/endpoint/__init__.py | 16 ++
.../tests/endpoint/test_data_lake_measure.py | 138 +++++++++++++
13 files changed, 607 insertions(+), 35 deletions(-)
diff --git a/streampipes-client-python/README.md
b/streampipes-client-python/README.md
index 67e94aaae..4231ae59c 100644
--- a/streampipes-client-python/README.md
+++ b/streampipes-client-python/README.md
@@ -72,7 +72,7 @@ pip install
git+https://github.com/apache/streampipes.git#subdirectory=streampip
>>> client.describe()
Hi there!
-You are connected to a StreamPipes instance running at http://localhost: 80.
+You are connected to a StreamPipes instance running at http://localhost:80.
The following StreamPipes resources are available with this client:
6x DataStreams
1x DataLakeMeasures
diff --git
a/streampipes-client-python/docs/examples/1-introduction-to-streampipes-python-client.ipynb
b/streampipes-client-python/docs/examples/1-introduction-to-streampipes-python-client.ipynb
index b034c02a7..9c6df65b9 100644
---
a/streampipes-client-python/docs/examples/1-introduction-to-streampipes-python-client.ipynb
+++
b/streampipes-client-python/docs/examples/1-introduction-to-streampipes-python-client.ipynb
@@ -3,11 +3,11 @@
{
"cell_type": "markdown",
"source": [
- "# Introduction to StreamPipes Python Client\n",
+ "# Introduction to StreamPipes Python\n",
"\n",
"<br>\n",
"\n",
- "### Why there is an extra Python client for StreamPipes\n",
+ "### Why there is an extra Python library for StreamPipes?\n",
"[Apache StreamPipes](https://streampipes.apache.org/) aims to enable
non-technical users to connect and analyze IoT data streams.\n",
"To this end, it provides an easy-to-use and convenient user interface
that allows one to connect to an IoT data source and create some visual\n",
"graphs within a few minutes. <br>\n",
@@ -20,10 +20,7 @@
"\n",
"<br>\n",
"\n",
- "### How to install the Python client\n",
- "Up to this point, we do not provide a release of the Python client in any
of the package indexes known for Python.\n",
- "This will probably start with StreamPipes `1.0.0` when we officially
launch the Python client.\n",
- "Until then, you can just install the currently available development
version of the client directly from GitHub.\n",
+ "### How to install StreamPipes Python?\n",
"Simply use the following `pip` command:"
],
"metadata": {
@@ -35,7 +32,17 @@
"execution_count": null,
"outputs": [],
"source": [
- "%pip install streampipes\n",
+ "%pip install streampipes\n"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
+ {
+ "cell_type": "code",
+ "execution_count": null,
+ "outputs": [],
+ "source": [
"# if you want to have the current development state you can also
execute\n",
"%pip install
git+https://github.com/apache/streampipes.git#subdirectory=streampipes-client-python"
],
@@ -43,6 +50,20 @@
"collapsed": false
}
},
+ {
+ "cell_type": "markdown",
+ "source": [
+ "### How to prepare the tutorials\n",
+ "In case you want to reproduce the first two tutorials exactly on your
end, you need to create a simple pipeline in StreamPipes like demonstrated
below.\n",
+ "\n",
+
"\n",
+ "\n",
+ "<br>"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
{
"cell_type": "markdown",
"source": [
@@ -62,7 +83,7 @@
},
{
"cell_type": "code",
- "execution_count": 3,
+ "execution_count": null,
"outputs": [],
"source": [
"from streampipes.client import StreamPipesClient\n",
@@ -81,7 +102,7 @@
"config = StreamPipesClientConfig(\n",
" credential_provider=StreamPipesApiKeyCredentials(\n",
" username=\"[email protected]\",\n",
- " api_key=\"DEMO-KEY\",\n",
+ " api_key=\"API-KEY\",\n",
" ),\n",
" host_address=\"localhost\",\n",
" https_disabled=True,\n",
@@ -121,7 +142,27 @@
"source": [
"Please note that you pass the names of the environment variables.\n",
"To ensure that the above code works, you must set the environment
variables with the same name you specified in `from_env`.\n",
- "\n",
+ "In this scenario this would look like the following:"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
+ {
+ "cell_type": "code",
+ "execution_count": null,
+ "outputs": [],
+ "source": [
+ "%export USER=\"<USERNAME>\"\n",
+ "%export"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
+ {
+ "cell_type": "markdown",
+ "source": [
"Having the `config` ready, we can now initialize the actual client."
],
"metadata": {
@@ -150,8 +191,23 @@
},
{
"cell_type": "code",
- "execution_count": null,
- "outputs": [],
+ "execution_count": 6,
+ "outputs": [
+ {
+ "name": "stdout",
+ "output_type": "stream",
+ "text": [
+ "2023-02-24 17:05:49,398 - streampipes.endpoint.endpoint - [INFO] -
[endpoint.py:167] [_make_request] - Successfully retrieved all resources.\n",
+ "2023-02-24 17:05:49,457 - streampipes.endpoint.endpoint - [INFO] -
[endpoint.py:167] [_make_request] - Successfully retrieved all resources.\n",
+ "\n",
+ "Hi there!\n",
+ "You are connected to a StreamPipes instance running at
http://localhost:80.\n",
+ "The following StreamPipes resources are available with this client:\n",
+ "1x DataLakeMeasures\n",
+ "1x DataStreams\n"
+ ]
+ }
+ ],
"source": [
"client.describe()"
],
diff --git
a/streampipes-client-python/docs/examples/2-extracting-data-from-the-streampipes-data-lake.ipynb
b/streampipes-client-python/docs/examples/2-extracting-data-from-the-streampipes-data-lake.ipynb
index b7001a6d8..780831c57 100644
---
a/streampipes-client-python/docs/examples/2-extracting-data-from-the-streampipes-data-lake.ipynb
+++
b/streampipes-client-python/docs/examples/2-extracting-data-from-the-streampipes-data-lake.ipynb
@@ -27,6 +27,19 @@
"collapsed": false
}
},
+ {
+ "cell_type": "code",
+ "execution_count": null,
+ "outputs": [],
+ "source": [
+ "# if you want all necessary dependencies required for this tutorial to be
installed,\n",
+ "# you can simply execute the following command\n",
+ "%pip install matplotlib streampipes"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
{
"cell_type": "code",
"execution_count": 2,
@@ -64,7 +77,7 @@
"name": "stdout",
"output_type": "stream",
"text": [
- "2022-12-04 21:19:21,832 - streampipes.client.client - [INFO] -
[client.py:127] [_set_up_logging] - Logging successfully initialized with
logging level INFO.\n"
+ "2023-02-24 17:34:25,860 - streampipes.client.client - [INFO] -
[client.py:128] [_set_up_logging] - Logging successfully initialized with
logging level INFO.\n"
]
}
],
@@ -94,7 +107,7 @@
"name": "stdout",
"output_type": "stream",
"text": [
- "2022-12-04 21:19:23,599 - streampipes.endpoint.endpoint - [INFO] -
[endpoint.py:153] [_make_request] - Successfully retrieved all resources.\n"
+ "2023-02-24 17:34:25,929 - streampipes.endpoint.endpoint - [INFO] -
[endpoint.py:167] [_make_request] - Successfully retrieved all resources.\n"
]
}
],
@@ -149,7 +162,7 @@
"outputs": [
{
"data": {
- "text/plain":
"DataLakeMeasure(element_id='urn:streampipes.apache.org:spi:datalakemeasure:xLSfXZ',
measure_name='test', timestamp_field='s0::timestamp',
event_schema=EventSchema(element_id='urn:streampipes.apache.org:spi:eventschema:UDMHXn',
event_properties=[EventProperty(element_id='urn:streampipes.apache.org:spi:eventpropertyprimitive:utvSWg',
label='Density', description='Denotes the current density of the fluid',
runtime_name='density', required=False, domain_properties=['http [...]
+ "text/plain":
"DataLakeMeasure(element_id='3cb6b5e6f107452483d1fd2ccf4bf9f9',
measure_name='test', timestamp_field='s0::timestamp',
event_schema=EventSchema(event_properties=[EventProperty(class_name='org.apache.streampipes.model.schema.EventPropertyPrimitive',
element_id='sp:eventproperty:EiFnkL', label='Density', description='Denotes
the current density of the fluid', runtime_name='density', required=False,
domain_properties=['http://schema.org/Number'], property_scope='MEASUREME [...]
},
"execution_count": 7,
"metadata": {},
@@ -178,8 +191,8 @@
"outputs": [
{
"data": {
- "text/plain": " measure_name timestamp_field pipeline_id pipeline_name
pipeline_is_running \\\n0 flow-rate s0::timestamp None
None False \n1 test s0::timestamp None
None False \n\n num_event_properties \n0
3 \n1 6 ",
- "text/html": "<div>\n<style scoped>\n .dataframe tbody tr
th:only-of-type {\n vertical-align: middle;\n }\n\n .dataframe
tbody tr th {\n vertical-align: top;\n }\n\n .dataframe thead th
{\n text-align: right;\n }\n</style>\n<table border=\"1\"
class=\"dataframe\">\n <thead>\n <tr style=\"text-align: right;\">\n
<th></th>\n <th>measure_name</th>\n <th>timestamp_field</th>\n
<th>pipeline_id</th>\n <th>pipeline_name</ [...]
+ "text/plain": " measure_name timestamp_field pipeline_id pipeline_name
pipeline_is_running \\\n0 flow-rate s0::timestamp None
None False \n1 test s0::timestamp None
None False \n\n num_event_properties \n0
6 \n1 6 ",
+ "text/html": "<div>\n<style scoped>\n .dataframe tbody tr
th:only-of-type {\n vertical-align: middle;\n }\n\n .dataframe
tbody tr th {\n vertical-align: top;\n }\n\n .dataframe thead th
{\n text-align: right;\n }\n</style>\n<table border=\"1\"
class=\"dataframe\">\n <thead>\n <tr style=\"text-align: right;\">\n
<th></th>\n <th>measure_name</th>\n <th>timestamp_field</th>\n
<th>pipeline_id</th>\n <th>pipeline_name</ [...]
},
"metadata": {},
"output_type": "display_data"
@@ -212,7 +225,7 @@
"name": "stdout",
"output_type": "stream",
"text": [
- "2022-12-04 21:19:30,505 - streampipes.endpoint.endpoint - [INFO] -
[endpoint.py:153] [_make_request] - Successfully retrieved all resources.\n"
+ "2023-02-24 17:34:26,020 - streampipes.endpoint.endpoint - [INFO] -
[endpoint.py:167] [_make_request] - Successfully retrieved all resources.\n"
]
}
],
@@ -258,7 +271,7 @@
"outputs": [
{
"data": {
- "text/plain": "2020"
+ "text/plain": "1000"
},
"execution_count": 11,
"metadata": {},
@@ -287,8 +300,8 @@
"outputs": [
{
"data": {
- "text/plain": " mass_flow temperature\ncount 2020.000000
2020.000000\nmean 4.976635 52.688616\nstd 2.920448
8.756244\nmin 0.003300 40.002800\n25% 2.443325 45.250551\n50%
4.886400 50.289900\n75% 7.524550 60.050674\nmax
9.997400 69.993896",
- "text/html": "<div>\n<style scoped>\n .dataframe tbody tr
th:only-of-type {\n vertical-align: middle;\n }\n\n .dataframe
tbody tr th {\n vertical-align: top;\n }\n\n .dataframe thead th
{\n text-align: right;\n }\n</style>\n<table border=\"1\"
class=\"dataframe\">\n <thead>\n <tr style=\"text-align: right;\">\n
<th></th>\n <th>mass_flow</th>\n <th>temperature</th>\n </tr>\n
</thead>\n <tbody>\n <tr>\n <th>count< [...]
+ "text/plain": " density mass_flow temperature
volume_flow\ncount 1000.000000 1000.000000 1000.000000 1000.000000\nmean
45.560337 5.457014 45.480231 5.659558\nstd 3.201544
3.184959 3.132878 3.122437\nmin 40.007698 0.004867
40.000992 0.039422\n25% 42.819497 2.654101 42.754623
3.021625\n50% 45.679264 5.382355 45.435944 5.572553\n75%
48.206881 8.183144 48.248473 8.338209\ [...]
+ "text/html": "<div>\n<style scoped>\n .dataframe tbody tr
th:only-of-type {\n vertical-align: middle;\n }\n\n .dataframe
tbody tr th {\n vertical-align: top;\n }\n\n .dataframe thead th
{\n text-align: right;\n }\n</style>\n<table border=\"1\"
class=\"dataframe\">\n <thead>\n <tr style=\"text-align: right;\">\n
<th></th>\n <th>density</th>\n <th>mass_flow</th>\n
<th>temperature</th>\n <th>volume_flow</th>\n </tr [...]
},
"execution_count": 12,
"metadata": {},
@@ -318,7 +331,7 @@
{
"data": {
"text/plain": "<Figure size 640x480 with 1 Axes>",
- "image/png":
"iVBORw0KGgoAAAANSUhEUgAAAh8AAAGdCAYAAACyzRGfAAAAOXRFWHRTb2Z0d2FyZQBNYXRwbG90bGliIHZlcnNpb24zLjYuMiwgaHR0cHM6Ly9tYXRwbG90bGliLm9yZy8o6BhiAAAACXBIWXMAAA9hAAAPYQGoP6dpAAC8Z0lEQVR4nOydd3gVRffHvzc9gRQIkNBB6b0oEFBBRIHXDvaO2AEFLMjvtYG+4mvF3l4FG6JYUEBAQHrvvYZOCjUJBFLv/v7Y3JvdvbOzM3v33tzA+TwPD7m7szOzu7MzZ845c8alKIoCgiAIgiCIIBFW0RUgCIIgCOLCgoQPgiAIgiCCCgkfBEEQBEEEFRI+CIIgCIIIKiR8EARBEAQRVEj4IAiCIAgiqJDwQRAEQRBEUCHhgyAIgiCIoBJR0RUw4na7kZGRgfj4eLhcroquDkEQBEEQAiiKgtOnT6NOnToIC+Pr
[...]
+ "image/png":
"iVBORw0KGgoAAAANSUhEUgAAAh8AAAGdCAYAAACyzRGfAAAAOXRFWHRTb2Z0d2FyZQBNYXRwbG90bGliIHZlcnNpb24zLjcuMCwgaHR0cHM6Ly9tYXRwbG90bGliLm9yZy88F64QAAAACXBIWXMAAA9hAAAPYQGoP6dpAACouklEQVR4nO2dd5jVxNfHv/duX9hCXdrSuxQpAgsqqCj62hDsqIBdAQWs/OwV7L2iggVEUUQRAREBKUvvvXd2qdth2837Rzb3TnInySQ3N7vA+TwPD3tTJpPJlDPnnDnjkSRJAkEQBEEQhEt4yzsDBEEQBEGcW5DwQRAEQRCEq5DwQRAEQRCEq5DwQRAEQRCEq5DwQRAEQRCEq5DwQRAEQRCEq5DwQRAEQRCEq5DwQRAEQRCEq0SWdwa0+Hw+HDp0CAkJCfB4POWdHYIgCIIgBJAkCbm5uahTpw68XmPdRoUTPg4d
[...]
},
"metadata": {},
"output_type": "display_data"
@@ -333,6 +346,126 @@
"collapsed": false
}
},
+ {
+ "cell_type": "markdown",
+ "source": [
+ "For data lake measurements, the `get()` method is even more powerful than
simply returning all the data for a given data lake measurement. We will look
at a selection of these below. The full list of supported parameters can be
found in the [docs](). <br>\n",
+ "Let's start by referring to the graph we created above, where we use only
two columns of our data lake measurement. If we already know this, we can
directly restrict the queried data to a subset of columns by using the
`columns` parameter. <br>\n",
+ "`columns` takes a list of column names as a comma-separated string:"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
+ {
+ "cell_type": "code",
+ "execution_count": 14,
+ "outputs": [
+ {
+ "name": "stdout",
+ "output_type": "stream",
+ "text": [
+ "2023-02-24 17:34:26,492 - streampipes.endpoint.endpoint - [INFO] -
[endpoint.py:167] [_make_request] - Successfully retrieved all resources.\n"
+ ]
+ },
+ {
+ "data": {
+ "text/plain": " time mass_flow temperature\n0
2023-02-24T16:19:41.472Z 3.309556 44.448483\n1
2023-02-24T16:19:41.482Z 5.608580 40.322033\n2 2023-02-24T16:19:41.493Z
7.692881 49.239639\n3 2023-02-24T16:19:41.503Z 3.632898
49.933754\n4 2023-02-24T16:19:41.513Z 0.711260 50.106617\n..
... ... ...\n995 2023-02-24T16:19:52.927Z
1.740114 46.558231\n996 2023-02-24T16:19:52.94Z [...]
+ "text/html": "<div>\n<style scoped>\n .dataframe tbody tr
th:only-of-type {\n vertical-align: middle;\n }\n\n .dataframe
tbody tr th {\n vertical-align: top;\n }\n\n .dataframe thead th
{\n text-align: right;\n }\n</style>\n<table border=\"1\"
class=\"dataframe\">\n <thead>\n <tr style=\"text-align: right;\">\n
<th></th>\n <th>time</th>\n <th>mass_flow</th>\n
<th>temperature</th>\n </tr>\n </thead>\n <tbody>\n < [...]
+ },
+ "execution_count": 14,
+ "metadata": {},
+ "output_type": "execute_result"
+ }
+ ],
+ "source": [
+ "flow_rate_pd = client.dataLakeMeasureApi.get(identifier=\"flow-rate\",
columns=\"mass_flow,temperature\").to_pandas()\n",
+ "flow_rate_pd"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
+ {
+ "cell_type": "markdown",
+ "source": [
+ "By default, the client returns only the first one thousand records of a
Data Lake measurement. This can be changed by passing a concrete value for the
`limit` parameter:"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
+ {
+ "cell_type": "code",
+ "execution_count": 15,
+ "outputs": [
+ {
+ "name": "stdout",
+ "output_type": "stream",
+ "text": [
+ "2023-02-24 17:34:26,736 - streampipes.endpoint.endpoint - [INFO] -
[endpoint.py:167] [_make_request] - Successfully retrieved all resources.\n"
+ ]
+ },
+ {
+ "data": {
+ "text/plain": "9528"
+ },
+ "execution_count": 15,
+ "metadata": {},
+ "output_type": "execute_result"
+ }
+ ],
+ "source": [
+ "flow_rate_pd = client.dataLakeMeasureApi.get(identifier=\"flow-rate\",
limit=10000).to_pandas()\n",
+ "len(flow_rate_pd)"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
+ {
+ "cell_type": "markdown",
+ "source": [
+ "If you want your data to be selected by time of occurrence rather than
quantity, you can specify your time window by passing the `start_date` and
`end_date` parameters:"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
+ {
+ "cell_type": "code",
+ "execution_count": 16,
+ "outputs": [
+ {
+ "name": "stdout",
+ "output_type": "stream",
+ "text": [
+ "2023-02-24 17:34:26,899 - streampipes.endpoint.endpoint - [INFO] -
[endpoint.py:167] [_make_request] - Successfully retrieved all resources.\n"
+ ]
+ },
+ {
+ "data": {
+ "text/plain": "<Figure size 640x480 with 1 Axes>",
+ "image/png":
"iVBORw0KGgoAAAANSUhEUgAAAh8AAAGdCAYAAACyzRGfAAAAOXRFWHRTb2Z0d2FyZQBNYXRwbG90bGliIHZlcnNpb24zLjcuMCwgaHR0cHM6Ly9tYXRwbG90bGliLm9yZy88F64QAAAACXBIWXMAAA9hAAAPYQGoP6dpAACdPUlEQVR4nO2dd3hb5fXHv1eSJe89k9ixs/dySOKEkACBEEaBhL1TKBTCCBRa8ivQRRvasltGSwuBFhr23oQkkL134kzHTrzteNuyLd3fH6/ee6/kK+lebcfn8zx+bEuy9Frj3u97zvecI4iiKIIgCIIgCCJEGMK9AIIgCIIg+hYkPgiCIAiCCCkkPgiCIAiCCCkkPgiCIAiCCCkkPgiCIAiCCCkkPgiCIAiCCCkkPgiCIAiCCCkkPgiCIAiCCCmmcC/AFbvdjvLyciQkJEAQhHAvhyAIgiAIDYiiiObmZvTr1w8G
[...]
+ },
+ "metadata": {},
+ "output_type": "display_data"
+ }
+ ],
+ "source": [
+ "from datetime import datetime\n",
+ "flow_rate_pd = client.dataLakeMeasureApi.get(\n",
+ " identifier=\"flow-rate\",\n",
+ " start_date=datetime(year=2023, month=2, day=24, hour=17, minute=21,
second=0),\n",
+ " end_date=datetime(year=2023, month=2, day=24, hour=17, minute=21,
second=1),\n",
+ " ).to_pandas()\n",
+ "flow_rate_pd.plot(y=[\"mass_flow\", \"temperature\"])\n",
+ "plt.show()"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
{
"cell_type": "markdown",
"source": [
diff --git
a/streampipes-client-python/docs/examples/3-getting-live-data-from-the-streampipes-data-stream.ipynb
b/streampipes-client-python/docs/examples/3-getting-live-data-from-the-streampipes-data-stream.ipynb
index b8eeb0d6b..8a9aba39b 100644
---
a/streampipes-client-python/docs/examples/3-getting-live-data-from-the-streampipes-data-stream.ipynb
+++
b/streampipes-client-python/docs/examples/3-getting-live-data-from-the-streampipes-data-stream.ipynb
@@ -35,6 +35,18 @@
"from streampipes.client.credential_provider import
StreamPipesApiKeyCredentials"
]
},
+ {
+ "cell_type": "code",
+ "execution_count": null,
+ "outputs": [],
+ "source": [
+ "# You can install all required libraries for this tutorial with the
following command\n",
+ "%pip install matplotlib ipython streampipes"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
{
"cell_type": "code",
"execution_count": 2,
diff --git
a/streampipes-client-python/docs/examples/4-using-online-machine-learning-on-a-streampipes-data-stream.ipynb
b/streampipes-client-python/docs/examples/4-using-online-machine-learning-on-a-streampipes-data-stream.ipynb
index bd546f17e..2e82de279 100644
---
a/streampipes-client-python/docs/examples/4-using-online-machine-learning-on-a-streampipes-data-stream.ipynb
+++
b/streampipes-client-python/docs/examples/4-using-online-machine-learning-on-a-streampipes-data-stream.ipynb
@@ -20,6 +20,18 @@
"from streampipes.client.credential_provider import
StreamPipesApiKeyCredentials"
]
},
+ {
+ "cell_type": "code",
+ "execution_count": null,
+ "outputs": [],
+ "source": [
+ "# you can install all required dependecies for this tutorial by executing
the following command\n",
+ "%pip install river streampipes"
+ ],
+ "metadata": {
+ "collapsed": false
+ }
+ },
{
"cell_type": "code",
"execution_count": 2,
@@ -312,7 +324,6 @@
"metadata": {},
"outputs": [],
"source": [
- "import pickle\n",
"from river import cluster, compose, preprocessing, tree\n",
"from streampipes.function_zoo.river_function import OnlineML\n",
"from streampipes.functions.utils.data_stream_generator import
RuntimeType\n",
diff --git a/streampipes-client-python/docs/img/tutorial-preparation.gif
b/streampipes-client-python/docs/img/tutorial-preparation.gif
new file mode 100644
index 000000000..863289c1d
Binary files /dev/null and
b/streampipes-client-python/docs/img/tutorial-preparation.gif differ
diff --git a/streampipes-client-python/streampipes/client/client.py
b/streampipes-client-python/streampipes/client/client.py
index 2c2c82dce..471696c1d 100644
--- a/streampipes-client-python/streampipes/client/client.py
+++ b/streampipes-client-python/streampipes/client/client.py
@@ -56,7 +56,7 @@ class StreamPipesClient:
--------
>>> from streampipes.client import StreamPipesClient
- >>> from streampipes.client.client_config import StreamPipesClientConfig
+ >>> from streampipes.client.config import StreamPipesClientConfig
>>> from streampipes.client.credential_provider import
StreamPipesApiKeyCredentials
>>> client_config = StreamPipesClientConfig(
diff --git
a/streampipes-client-python/streampipes/endpoint/api/data_lake_measure.py
b/streampipes-client-python/streampipes/endpoint/api/data_lake_measure.py
index b39ca5d2b..a5975b277 100644
--- a/streampipes-client-python/streampipes/endpoint/api/data_lake_measure.py
+++ b/streampipes-client-python/streampipes/endpoint/api/data_lake_measure.py
@@ -19,8 +19,10 @@
Specific implementation of the StreamPipes API's data lake measure endpoints.
This endpoint allows to consume data stored in StreamPipes' data lake
"""
-from typing import Tuple, Type
+from datetime import datetime
+from typing import Any, Dict, Literal, Optional, Tuple, Type
+from pydantic import BaseModel, Extra, Field, StrictInt, ValidationError,
validator
from streampipes.endpoint.endpoint import APIEndpoint
from streampipes.model.container import DataLakeMeasures
from streampipes.model.container.resource_container import ResourceContainer
@@ -31,8 +33,130 @@ __all__ = [
]
+class StreamPipesQueryValidationError(Exception):
+ """A custom exception to be raised when the validation of query parameter
+ causes an error.
+ """
+
+
+class MeasurementGetQueryConfig(BaseModel):
+ """Config class describing the parameters of the GET endpoint for
measurements.
+
+ This config class is used to validate the provided query parameters for
the GET endpoint of measurements.
+ Additionally, it takes care of the conversion to a proper HTTP query
string.
+ Thereby, parameter names are adapted to the naming of the StreamPipes API,
for which Pydantic aliases are used.
+
+ Attributes
+ ----------
+ columns: Optional[str]
+ A comma separated list of column names (e.g., `time,value`)<br>
+ If provided, the returned data only consists of the given columns.<br>
+ Please be aware that the column `time` as an index is always included.
+ end_date: Optional[datetime]
+ Restricts queried data to be younger than the specified time.
+ limit: Optional[int]
+ Amount of records returned at maximum (default: `1000`) <br>
+ This needs to be at least `1`
+ offset: Optional[int]
+ Offset to be applied to returned data <br>
+ This needs to be at least `0`
+ order: Optional[str]
+ Ordering of query results <br>
+ Allowed values: `ASC` and `DESC` (default: `ASC`)
+ page_no: Optional[int]
+ Page number used for paging operation <br>
+ This needs to be at least `1`
+ start_date: Optional[datetime]
+ Restricts queried data to be older than the specified time
+ """
+
+ _regex_comma_separated_string = r"^[0-9a-zA-Z\_]+(,[0-9a-zA-Z\_]+)*$"
+
+ class Config:
+ """Pydantic Config class"""
+
+ extra = Extra.forbid
+ allow_population_by_field_name = True
+
+ columns: Optional[str] = Field(regex=_regex_comma_separated_string)
+ end_date: Optional[StrictInt] = Field(alias="endDate")
+ limit: Optional[int] = Field(ge=1, default=1000)
+ offset: Optional[int] = Field(ge=0)
+ order: Optional[Literal["ASC", "DESC"]]
+ page_no: Optional[int] = Field(alias="page", ge=1)
+ start_date: Optional[StrictInt] = Field(alias="startDate")
+
+ @validator("end_date", "start_date", pre=True)
+ @classmethod
+ def _convert_datetime(cls, dt: datetime) -> int:
+ """Pydantic validator to convert datetime object to unix timestamp.
+
+ The StreamPipes API expects datetime related parameters to be passed
as unix timestamp.
+ For the sake of convenience we expect datetime objects to be passed
for these values.
+ This requires us to convert the provided datetime objects in unix
timestamp representation
+
+ Parameters
+ ----------
+ dt: datetime
+ The datetime value to be passed as query parameter
+
+
+ Raises
+ ------
+ StreamPipesQueryValidationError
+ In case `start_date` or `end_date` is not passed as a datetime
object
+ ValueError
+ In case the transformation of the datetime object did not work
+
+ Returns
+ -------
+ unix_timestamp: int
+ unix timestamp of the given timestamp
+
+ """
+
+ if not isinstance(dt, datetime) or dt is None:
+ raise StreamPipesQueryValidationError(
+ f"The passed value for either `start_date` or `end_date` "
f"is not a datetime object: '{dt}'."
+ )
+ try:
+ unix_timestamp = int(datetime.timestamp(dt) * 1000)
+ return unix_timestamp
+ except ValueError as ve: # pragma: no cover
+ raise ValueError(
+ "Your datetime object is off, it could not be parsed"
+ "This should not occur, but unfortunately did.\n"
+ "Therefore, it would be great if you could report this problem
as an issue at "
+ "github.com/apache/streampipes.\n"
+ ) from ve
+
+ def build_query_string(self) -> str:
+ """Builds a HTTP query string for the config.
+
+ This method returns an HTTP query string for the invoking config.
+ It follows the following structure `?param1=value1¶m2=value2...`.
+ This query string is not an entire URL, instead it needs to appended
to an API path.
+
+ Returns
+ -------
+ query_param_string: str
+ HTTP query params string (`?param1=value1¶m2=value2...`)
+ """
+
+ # create dictionary representation of the config that meets the
following expectations:
+ # - query parameter should comply to the parameter names of the
StreamPipes API (`by_alias`)
+ # - query params should only be present if they are different from
None (`exclude_none`)
+ query_param_dict = self.dict(by_alias=True, exclude_none=True)
+
+ # create query string that complies to HTTP syntax
(?param1=value1¶m2=value2&...)
+ query_param_string = f"?{'&'.join([f'{k}={v}' for k, v in
query_param_dict.items()])}"
+
+ return query_param_string
+
+
class DataLakeMeasureEndpoint(APIEndpoint):
"""Implementation of the DataLakeMeasure endpoint.
+
This endpoint provides an interfact to all data stored in the StreamPipes
data lake.
Consequently, it allows uerying metadata about available data sets (see
`all()` method).
@@ -50,7 +174,7 @@ class DataLakeMeasureEndpoint(APIEndpoint):
--------
>>> from streampipes.client import StreamPipesClient
- >>> from streampipes.client.client_config import StreamPipesClientConfig
+ >>> from streampipes.client.config import StreamPipesClientConfig
>>> from streampipes.client.credential_provider import
StreamPipesApiKeyCredentials
>>> client_config = StreamPipesClientConfig(
@@ -66,8 +190,72 @@ class DataLakeMeasureEndpoint(APIEndpoint):
>>> len(data_lake_measures)
5
+
+ Retrieve a specific data lake measure as a pandas DataFrame
+ >>> flow_rate_pd =
client.dataLakeMeasureApi.get(identifier="flow-rate").to_pandas()
+ >>> flow_rate_pd
+ time density mass_flow sensorId
sensor_fault_flags temperature volume_flow
+ 0 2023-02-24T16:19:41.472Z 50.872730 3.309556 flowrate02
False 44.448483 5.793138
+ 1 2023-02-24T16:19:41.482Z 47.186588 5.608580 flowrate02
False 40.322033 0.058015
+ 2 2023-02-24T16:19:41.493Z 46.735321 7.692881 flowrate02
False 49.239639 10.283526
+ 3 2023-02-24T16:19:41.503Z 40.169796 3.632898 flowrate02
False 49.933754 6.893441
+ 4 2023-02-24T16:19:41.513Z 49.635124 0.711260 flowrate02
False 50.106617 2.999871
+ .. ... ... ... ...
... ... ...
+ 995 2023-02-24T16:19:52.927Z 50.057495 1.740114 flowrate02
False 46.558231 1.818237
+ 996 2023-02-24T16:19:52.94Z 41.038895 7.211723 flowrate02
False 48.048622 2.127493
+ 997 2023-02-24T16:19:52.952Z 45.837013 7.770180 flowrate02
False 48.188026 7.892062
+ 998 2023-02-24T16:19:52.965Z 43.389065 4.458602 flowrate02
False 48.280899 5.733892
+ 999 2023-02-24T16:19:52.977Z 44.056030 2.592060 flowrate02
False 47.505951 4.260697
+
+ As you can see, the returned amount of rows per default is `1000`.
+ We can modify this behavior by passing the `limit` paramter.
+ >>> flow_rate_pd = client.dataLakeMeasureApi.get(identifier="flow-rate",
limit=10).to_pandas()
+ >>> len(flow_rate_pd)
+
+ If we are only interested in the values for `density`,
+ `columns` allows us to select the columns to be returned:
+ >>> flow_rate_pd = client.dataLakeMeasureApi.get(identifier="flow-rate",
columns='density', limit=3).to_pandas()
+ >>> flow_rate_pd
+ time density
+ 0 2023-02-24T16:19:41.472Z 50.872730
+ 1 2023-02-24T16:19:41.482Z 47.186588
+ 2 2023-02-24T16:19:41.493Z 46.735321
+
+ This is only a subset of the available query parameters,
+ find them at
[MeasurementGetQueryConfig][streampipes.endpoint.api.data_lake_measure.MeasurementGetQueryConfig].
"""
+ @staticmethod
+ def _validate_query_params(query_params: Dict[str, Any]) ->
MeasurementGetQueryConfig:
+ """Validates given query params.
+
+ Validates the given query parameters via the
+
[MeasurementGetQueryConfig][streampipes.endpoint.api.data_lake_measure.MeasurementGetQueryConfig].
+
+ Raises
+ ------
+ StreamPipesQueryValidationError
+ In case the query parameters are not provided correctly
+
+ Returns
+ -------
+ config: MeasurementGetQueryConfig
+ validated config that can be used to construct the query
+ """
+ try:
+ config = MeasurementGetQueryConfig.parse_obj(query_params)
+ except ValidationError as ve:
+ raise StreamPipesQueryValidationError(
+ f"\nOops, there seems to be a problem with your provided query
options. "
+ f"Some of them are not provided as expected. Please see the
detailed output below:\n\n"
+ f"Validation error log: {ve.json()}\n\n"
+ f"In case you assess your query configuration to be correct
feel free to file us an issue via "
+ f"github.com/apache/streampipes.\n"
+ f"Please don't forget to include the following validation log
from above."
+ )
+
+ return config
+
@property
def _resource_cls(self) -> Type[DataLakeSeries]:
"""
@@ -102,20 +290,38 @@ class DataLakeMeasureEndpoint(APIEndpoint):
return "api", "v4", "datalake", "measurements"
- def get(self, identifier: str) -> DataLakeSeries:
+ def get(self, identifier: str, **kwargs: Optional[Dict[str, Any]]) ->
DataLakeSeries:
"""Queries the specified data lake measure from the API.
+ By default, the maximum number of returned records is 1000.
+ This behaviour can be influences by passing the parameter `limit` with
a different value
+ (see
[MeasurementGetQueryConfig][streampipes.endpoint.api.data_lake_measure.MeasurementGetQueryConfig]).
+
Parameters
----------
identifier: str
The identifier of the data lake measure to be queried.
+ **kwargs: Dict[str, Any]
+ keyword arguments can be used to provide additional query
parameters.
+ The available query parameters are defined by the
+
[MeasurementGetQueryConfig][streampipes.endpoint.api.data_lake_measure.MeasurementGetQueryConfig].
Returns
-------
- The specified data lake measure as an instance of the corresponding
model class (`model.DataLakeSeries`).
+ measurement: DataLakeMeasures
+ the specified data lake measure
+
+ Examples
+ --------
+ see directly at
[DataLakeMeasureEndpoint][streampipes.endpoint.api.data_lake_measure.DataLakeMeasureEndpoint].
"""
- response = self._make_request(
- request_method=self._parent_client.request_session.get,
url=f"{self.build_url()}/{identifier}"
- )
+ # bild base URL for resource
+ url = f"{self.build_url()}/{identifier}"
+
+ # extend base URL by query parameters
+ measurement_get_config =
self._validate_query_params(query_params=kwargs)
+ url += measurement_get_config.build_query_string()
+
+ response =
self._make_request(request_method=self._parent_client.request_session.get,
url=url)
return self._resource_cls.from_json(json_string=response.text)
diff --git a/streampipes-client-python/streampipes/endpoint/api/data_stream.py
b/streampipes-client-python/streampipes/endpoint/api/data_stream.py
index 71c77c969..4780d3ea4 100644
--- a/streampipes-client-python/streampipes/endpoint/api/data_stream.py
+++ b/streampipes-client-python/streampipes/endpoint/api/data_stream.py
@@ -45,7 +45,7 @@ class DataStreamEndpoint(APIEndpoint):
--------
>>> from streampipes.client import StreamPipesClient
- >>> from streampipes.client.client_config import StreamPipesClientConfig
+ >>> from streampipes.client.config import StreamPipesClientConfig
>>> from streampipes.client.credential_provider import
StreamPipesApiKeyCredentials
>>> client_config = StreamPipesClientConfig(
diff --git a/streampipes-client-python/streampipes/endpoint/endpoint.py
b/streampipes-client-python/streampipes/endpoint/endpoint.py
index b3ae4fe2d..cfc6b7381 100644
--- a/streampipes-client-python/streampipes/endpoint/endpoint.py
+++ b/streampipes-client-python/streampipes/endpoint/endpoint.py
@@ -193,7 +193,7 @@ class APIEndpoint(Endpoint):
)
return self._container_cls.from_json(json_string=response.text)
- def get(self, identifier: str) -> Resource:
+ def get(self, identifier: str, **kwargs) -> Resource:
"""Queries the specified resource from the API endpoint.
Parameters
diff --git a/streampipes-client-python/tests/client/test_data_lake_series.py
b/streampipes-client-python/tests/client/test_data_lake_series.py
index 49b75a90c..482c0f75b 100644
--- a/streampipes-client-python/tests/client/test_data_lake_series.py
+++ b/streampipes-client-python/tests/client/test_data_lake_series.py
@@ -115,7 +115,7 @@ class TestDataLakeSeries(TestCase):
result = client.dataLakeMeasureApi.get(identifier="test")
http_session.assert_has_calls(
-
[call().get(url="https://localhost:80/streampipes-backend/api/v4/datalake/measurements/test")],
+
[call().get(url="https://localhost:80/streampipes-backend/api/v4/datalake/measurements/test?limit=1000")],
any_order=True,
)
diff --git a/streampipes-client-python/tests/endpoint/__init__.py
b/streampipes-client-python/tests/endpoint/__init__.py
new file mode 100644
index 000000000..cce3acad3
--- /dev/null
+++ b/streampipes-client-python/tests/endpoint/__init__.py
@@ -0,0 +1,16 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
diff --git a/streampipes-client-python/tests/endpoint/test_data_lake_measure.py
b/streampipes-client-python/tests/endpoint/test_data_lake_measure.py
new file mode 100644
index 000000000..4ee22a607
--- /dev/null
+++ b/streampipes-client-python/tests/endpoint/test_data_lake_measure.py
@@ -0,0 +1,138 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+from datetime import datetime
+from unittest import TestCase
+from streampipes.endpoint.api.data_lake_measure import
DataLakeMeasureEndpoint, StreamPipesQueryValidationError
+
+
+class TestMeasurementGetQueryConfig(TestCase):
+
+ def test_default(self):
+ config_dict = {}
+ measurement_config =
DataLakeMeasureEndpoint._validate_query_params(query_params=config_dict)
+ result = measurement_config.build_query_string()
+
+ self.assertEqual("?limit=1000", result)
+
+ def test_additional_param_given(self):
+ config_dict = {
+ "columns": "time,value_25"
+ }
+
+ measurement_config =
DataLakeMeasureEndpoint._validate_query_params(query_params=config_dict)
+ result = measurement_config.build_query_string()
+
+ self.assertEqual("?columns=time,value_25&limit=1000", result)
+
+ def test_extra_param(self):
+ config_dict = {
+ "foo": "bar"
+ }
+
+ with self.assertRaises(StreamPipesQueryValidationError):
+
DataLakeMeasureEndpoint._validate_query_params(query_params=config_dict)
+
+ def test_alias_as_query_param(self):
+ config_dict = {
+ "page_no": 5
+ }
+
+ measurement_config =
DataLakeMeasureEndpoint._validate_query_params(query_params=config_dict)
+ result = measurement_config.build_query_string()
+
+ self.assertEqual("?limit=1000&page=5", result)
+
+ def test_datetime_validation(self):
+ now = datetime.utcnow()
+
+ config_dict = {
+ "start_date": now,
+ "end_date": now
+ }
+ measurement_config =
DataLakeMeasureEndpoint._validate_query_params(query_params=config_dict)
+ result = measurement_config.build_query_string()
+
+ expected_ts = int(datetime.timestamp(now) * 1000)
+ expected = f"?endDate={expected_ts}&limit=1000&startDate={expected_ts}"
+
+ self.assertEqual(expected, result)
+
+ def test_datetime_validation_no_datetime(self):
+ config_dict = {
+ "start_date": "test"
+ }
+
+ with self.assertRaises(StreamPipesQueryValidationError):
+
DataLakeMeasureEndpoint._validate_query_params(query_params=config_dict)
+
+ def test_columns_validation(self):
+ config_dict_one_col = {
+ "columns": "col1"
+ }
+ config_dict_mul_col = {
+ "columns": "col1,col2,col3"
+ }
+ config_dict_semicolon = {
+ "columns": "col1;col2"
+ }
+ config_dict_whitespace_ending = {
+ "columns": "col1 "
+ }
+
+ self.assertEqual("?columns=col1&limit=1000",
DataLakeMeasureEndpoint._validate_query_params(
+ query_params=config_dict_one_col).build_query_string())
+ self.assertEqual("?columns=col1,col2,col3&limit=1000",
DataLakeMeasureEndpoint._validate_query_params(
+ query_params=config_dict_mul_col).build_query_string())
+ with self.assertRaises(StreamPipesQueryValidationError):
+
DataLakeMeasureEndpoint._validate_query_params(query_params=config_dict_semicolon)
+ with self.assertRaises(StreamPipesQueryValidationError):
+
DataLakeMeasureEndpoint._validate_query_params(query_params=config_dict_whitespace_ending)
+
+ def test_minium_parameter_values(self):
+ config_dict_happy_path = {
+ "limit": 15,
+ "page_no": 3
+ }
+
+ config_dict_limit_too_low = {
+ "limit": 0
+ }
+
+ config_dict_page_no_too_low = {
+ "page_no": -2
+ }
+
+ measurement_config_happy =
DataLakeMeasureEndpoint._validate_query_params(query_params=config_dict_happy_path)
+ result_happy = measurement_config_happy.build_query_string()
+
+ self.assertEqual("?limit=15&page=3", result_happy)
+
+ with self.assertRaises(StreamPipesQueryValidationError):
+
DataLakeMeasureEndpoint._validate_query_params(query_params=config_dict_limit_too_low)
+
+ with self.assertRaises(StreamPipesQueryValidationError):
+
DataLakeMeasureEndpoint._validate_query_params(query_params=config_dict_page_no_too_low)
+
+ def test_literal_validation(self):
+
+ config_invalid_order = {
+ "order": "UP"
+ }
+
+ with self.assertRaises(StreamPipesQueryValidationError):
+
DataLakeMeasureEndpoint._validate_query_params(query_params=config_invalid_order)
\ No newline at end of file