daroczig / daroczig/AWR.Kinesis
Corrupted protocol when using consumer and producer in one session
- Dominant language
- R
- Stars
- 4
- Forks
- 1
- PR merge metrics
- No merged PRs in 30d
Description
It seems 'Execution halted' from R's stderr stream corrupts KCL protocol when using consumer and producer in one session.
I'm trying to implement `read_stream->transform->write_stream` pattern with `AWR.Kinesis`.
```r
library(futile.logger)
library(AWR.Kinesis)
flog.threshold(INFO)
kinesis_consumer(
processRecords = function(records) {
flog.info('Received %d records from Kinesis', nrow(records))
for(record in records$data)
kinesis_put_record("my-stream", record, "partition-1")
},
logfile = "log.txt"
)
```
At some point of time I see
> Jun 26, 2018 6:16:09 PM com.amazonaws.services.kinesis.multilang.DrainChildSTDERRTask handleLine
SEVERE: Received error line from subprocess [Execution halted] for shard shardId-000000000000
Execution halted
Followed by:
Click to expand
> Jun 26, 2018 6:16:09 PM com.amazonaws.services.kinesis.multilang.LineReaderTask call
INFO: Stopping: Reading STDERR for shardId-000000000000
Jun 26, 2018 6:16:27 PM com.amazonaws.services.kinesis.multilang.MessageWriter writeMessage
INFO: Writing ProcessRecordsMessage to child process for shard shardId-000000000000
Jun 26, 2018 6:16:27 PM com.amazonaws.services.kinesis.multilang.LineReaderTask call
INFO: Starting: Reading next message from STDIN for shardId-000000000000
Jun 26, 2018 6:16:27 PM com.amazonaws.services.kinesis.multilang.LineReaderTask call
INFO: Stopping: Reading next message from STDIN for shardId-000000000000
Jun 26, 2018 6:16:27 PM com.amazonaws.services.kinesis.multilang.MultiLangProtocol waitForStatusMessage
SEVERE: Failed to get status message for processRecords action for shard shardId-000000000000
java.util.concurrent.ExecutionException: java.lang.RuntimeException: Reached end of STDIN of child process for shard shardId-000000000000 so won't be able to return a message.
at java.util.concurrent.FutureTask.report(FutureTask.java:122)
at java.util.concurrent.FutureTask.get(FutureTask.java:192)
at com.amazonaws.services.kinesis.multilang.MultiLangProtocol.waitForStatusMessage(MultiLangProtocol.java:152)
at com.amazonaws.services.kinesis.multilang.MultiLangProtocol.waitForStatusMessage(MultiLangProtocol.java:120)
at com.amazonaws.services.kinesis.multilang.MultiLangProtocol.processRecords(MultiLangProtocol.java:86)
at com.amazonaws.services.kinesis.multilang.MultiLangRecordProcessor.processRecords(MultiLangRecordProcessor.java:100)
at com.amazonaws.services.kinesis.clientlibrary.lib.worker.ProcessTask.call(ProcessTask.java:176)
at com.amazonaws.services.kinesis.clientlibrary.lib.worker.MetricsCollectingTaskDecorator.call(MetricsCollectingTaskDecorator.java:49)
at com.amazonaws.services.kinesis.clientlibrary.lib.worker.MetricsCollectingTaskDecorator.call(MetricsCollectingTaskDecorator.java:24)
at java.util.concurrent.FutureTask.run(FutureTask.java:266)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)
Caused by: java.lang.RuntimeException: Reached end of STDIN of child process for shard shardId-000000000000 so won't be able to return a message.
at com.amazonaws.services.kinesis.multilang.GetNextMessageTask.returnAfterEndOfInput(GetNextMessageTask.java:84)
at com.amazonaws.services.kinesis.multilang.GetNextMessageTask.returnAfterEndOfInput(GetNextMessageTask.java:31)
at com.amazonaws.services.kinesis.multilang.LineReaderTask.call(LineReaderTask.java:70)
... 4 more
Jun 26, 2018 6:16:27 PM com.amazonaws.services.kinesis.multilang.MultiLangRecordProcessor stopProcessing
SEVERE: Encountered an error while trying to process records
java.lang.RuntimeException: Child process failed to process records
at com.amazonaws.services.kinesis.multilang.MultiLangRecordProcessor.processRecords(MultiLangRecordProcessor.java:101)
at com.amazonaws.services.kinesis.clientlibrary.lib.worker.ProcessTask.call(ProcessTask.java:176)
at com.amazonaws.services.kinesis.clientlibrary.lib.worker.MetricsCollectingTaskDecorator.call(MetricsCollectingTaskDecorator.java:49)
at com.amazonaws.services.kinesis.clientlibrary.lib.worker.MetricsCollectingTaskDecorator.call(MetricsCollectingTaskDecorator.java:24)
at java.util.concurrent.FutureTask.run(FutureTask.java:266)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)
Jun 26, 2018 6:16:27 PM com.amazonaws.services.kinesis.multilang.LineReaderTask call
INFO: Starting: Draining STDOUT for shardId-000000000000
Jun 26, 2018 6:16:27 PM com.amazonaws.services.kinesis.multilang.LineReaderTask call
INFO: Stopping: Draining STDOUT for shardId-000000000000
Jun 26, 2018 6:16:27 PM com.amazonaws.services.kinesis.multilang.MultiLangRecordProcessor childProcessShutdownSequence
INFO: Child process exited with value: 1
---------------
It looks like `kinesis_put_record` corrupts protocol (by creating another kinesis client in the same session?) Will it make sense to export more low-level interface to kinesis? I can send PR
Contributor guide
No contributing guide indexed for this repository
Research direction
Start by reproducing the read_stream-to-transform-to-write_stream example with AWR.Kinesis, focusing on the interaction between kinesis_consumer and kinesis_put_record in one session. Done means the consumer remains connected and the KCL protocol does not terminate with “Execution halted” or an end-of-STDIN error.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- aws, r
- Domain
- stream-processing
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 30/100