I have a dataframe that has columns with structs that have dates and values so that the schema looks like
root
|-- col1: struct (nullable = true)
| |-- dates: array (nullable = true)
| | |-- element: timestamp (containsNull = true)
| |-- values: array (nullable = true)
| | |-- element: double (containsNull = true)
|-- col2: struct (nullable = true)
| |-- dates: array (nullable = true)
| | |-- element: timestamp (containsNull = true)
| |-- values: array (nullable = true)
| | |-- element: double (containsNull = true)
|-- id: string (nullable = true)
Given some time index:
time_index = datetime.datetime(2015, 12, 12, 4, 45)
and the number of days before and after:
min_diff = -1 and max_diff = 2
I want new columns, col1_filt and col2_filt, with the same structure that return the dates that fall within the window defined by the time_index and min_diff and max_diff and the corresponding values. If none of the dates or values fall within that window I want it to return None.
Below is an example DataFrame to work with.
Example DataFrame:
example_input = [
Row(
id = "A",
col1 = Row(
dates = [datetime.datetime(2015, 12, 11, 5, 28), datetime.datetime(2015, 12, 12, 4, 45), datetime.datetime(2015, 12, 13, 5, 9)],
values = [17.7, 19.1, 19.1]
),
col2 = Row(
dates = [datetime.datetime(2015, 12, 13, 4, 48), datetime.datetime(2015, 12, 15, 5, 8)],
values = [19.1, 19.1]
)
),
Row(
id = "B",
col1 = Row(
dates = [datetime.datetime(2017, 1, 13, 5, 9)],
values = [19.1]
),
col2 = Row(
dates = [datetime.datetime(2017, 1, 12, 2, 48), datetime.datetime(2017, 1, 15, 5, 8)],
values = [19.5, 29.1]
)
),
]
df = spark.createDataFrame(example_input)
Showing df:
+-------------------------------------------------------------------------------------+----------------------------------------------------------+---+
|col1 |col2 |id |
+-------------------------------------------------------------------------------------+----------------------------------------------------------+---+
|[[2015-12-11 05:28:00, 2015-12-12 04:45:00, 2015-12-13 05:09:00], [17.7, 19.1, 19.1]]|[[2015-12-13 04:48:00, 2015-12-15 05:08:00], [19.1, 19.1]]|A |
|[[2017-01-13 05:09:00], [19.1]] |[[2017-01-12 02:48:00, 2017-01-15 05:08:00], [19.5, 29.1]]|B |
+-------------------------------------------------------------------------------------+----------------------------------------------------------+---+
I have some code that will take a Pyspark row object and return the filtered Pyspark row object but I am not sure how to make this a udf.