diff --git a/python/pyspark/streaming/dstream.py b/python/pyspark/streaming/dstream.py index 66024d539ce5c..a36f4b9bf9d87 100644 --- a/python/pyspark/streaming/dstream.py +++ b/python/pyspark/streaming/dstream.py @@ -436,6 +436,7 @@ def __init__(self, prev, func, preservesPartitioning=False): self._prev_jrdd_deserializer = prev._jrdd_deserializer else: prev_func = prev.func + def pipeline_func(split, iterator): return func(split, prev_func(split, iterator)) self.func = pipeline_func