Skip to content

Commit

Permalink
ML-345: null in TS column (#193)
Browse files Browse the repository at this point in the history
* ML-345: null in TS column

* pr comments

* minor fix
  • Loading branch information
katyakats authored Apr 6, 2021
1 parent be192b8 commit 9a7e12b
Show file tree
Hide file tree
Showing 2 changed files with 42 additions and 1 deletion.
34 changes: 34 additions & 0 deletions integration/test_flow_integration.py
Original file line number Diff line number Diff line change
Expand Up @@ -462,3 +462,37 @@ def test_write_multiple_keys_to_v3io_from_csv(setup_teardown_test):
assert response.status_code == 200
assert expected == response.output.item


def test_write_none_time(setup_teardown_test):

table = Table(setup_teardown_test, V3ioDriver())
data = pd.DataFrame(
{
"first_name": ["moshe", "yosi"],
"color": ['blue', 'yellow'],
"time": [test_base_time, None]
}
)

def set_moshe_time_to_none(data):
if data['first_name'] == 'moshe':
data['time'] = pd.NaT
return data

controller = build_flow([
DataframeSource(data, key_field='first_name'),
WriteToTable(table),
Map(set_moshe_time_to_none),
WriteToTable(table),
]).run()
controller.await_termination()

response = asyncio.run(get_kv_item(setup_teardown_test, 'yosi'))
expected = {'first_name': 'yosi', 'color': 'yellow'}
assert response.status_code == 200
assert expected == response.output.item

response = asyncio.run(get_kv_item(setup_teardown_test, 'moshe'))
expected = {'first_name': 'moshe', 'color': 'blue'}
assert response.status_code == 200
assert expected == response.output.item
9 changes: 8 additions & 1 deletion storey/drivers.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
import os
from datetime import datetime, timedelta
from typing import Optional
import pandas as pd

import v3io
import v3io.aio.dataplane
Expand Down Expand Up @@ -223,7 +224,11 @@ def _build_feature_store_update_expression(self, aggregation_element, additional
for name, value in additional_data.items():
if name.casefold() in self.saved_engine_words.keys():
name = f'`{name}`'
expressions.append(f'{name}={self._convert_python_obj_to_expression_value(value)}')
expression_value = self._convert_python_obj_to_expression_value(value)
if expression_value:
expressions.append(f'{name}={self._convert_python_obj_to_expression_value(value)}')
else:
expressions.append(f'REMOVE {name}')

update_expression = ';'.join(expressions)
return update_expression, condition_expression, pending_updates
Expand Down Expand Up @@ -352,6 +357,8 @@ def _convert_python_obj_to_expression_value(value):
elif isinstance(value, bytes):
return f"blob('{base64.b64encode(value).decode('ascii')}')"
elif isinstance(value, datetime):
if pd.isnull(value):
return None
timestamp = value.timestamp()

secs = int(timestamp)
Expand Down

0 comments on commit 9a7e12b

Please sign in to comment.