256a3a8013
## What changes were proposed in this pull request? This pr is to fix an issue occurred when resharding Kinesis streams; the resharding makes the KCL throw an exception because Spark does not checkpoint `SHARD_END` when finishing reading closed shards in `KinesisRecordProcessor#shutdown`. This bug finally leads to stopping subscribing new split (or merged) shards. ## How was this patch tested? Added a test in `KinesisStreamSuite` to check if it works well when splitting/merging shards. Author: Takeshi YAMAMURO <linguin.m.s@gmail.com> Closes #16213 from maropu/SPARK-18020. |
||
---|---|---|
.. | ||
__init__.py | ||
context.py | ||
dstream.py | ||
flume.py | ||
kafka.py | ||
kinesis.py | ||
listener.py | ||
tests.py | ||
util.py |