daroczig / daroczig/AWR.Kinesis

Corrupted protocol when using consumer and producer in one session

Open
#1 5 comments 0 reactions 0 assignees View on GitHub
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

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.