flink-dev mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From "Stephan Ewen (JIRA)" <j...@apache.org>
Subject [jira] [Created] (FLINK-5048) Kafka Consumer (0.9/0.10) threading model leads problematic cancellation behavior
Date Thu, 10 Nov 2016 10:30:59 GMT
Stephan Ewen created FLINK-5048:

             Summary: Kafka Consumer (0.9/0.10) threading model leads problematic cancellation
                 Key: FLINK-5048
                 URL: https://issues.apache.org/jira/browse/FLINK-5048
             Project: Flink
          Issue Type: Bug
          Components: Kafka Connector
    Affects Versions: 1.1.3
            Reporter: Stephan Ewen
            Assignee: Stephan Ewen
             Fix For: 1.2.0

The {{FLinkKafkaConsumer}} (0.9 / 0.10) spawns a separate thread that operates the KafkaConsumer.
That thread is shielded from interrupts, because the Kafka Consumer has not been handling
thread interrupts well.

Since that thread is also the thread that emits records, it may block in the network stack
(backpressure) or in chained operators. The later case leads to situations where cancellations
get very slow unless that thread would be interrupted (which it cannot be).

I propose to change the thread model as follows:

  - A spawned consumer thread pull from the KafkaConsumer and pushes its pulled batch of records
into a blocking queue (size one)
  - The main thread of the task will pull the record batches from the blocking queue and emit
the records.

This allows actually for some additional I/O overlay while limiting the additional memory
consumption - only two batches are ever held, one being fetched and one being emitted.

This message was sent by Atlassian JIRA

View raw message