DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
| Table of Contents |
|---|
Status
Current state: "Under Discussion"
...
Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).
Motivation
In transactions, having a poison record or sending the record to a topic with issues ends up going to an error state, which leads to dropping the whole batch. This logic does not allow the user to do advanced error handling and ignore some errors. More specifically, the custom production exception handler of Kafka Streams can not influence the control to skip the error by dropping the bad record and committing the transaction with the rest of the batch records.
...
We aim at enabling users to ignore the `send(ProducerRecord)` errors by explicitly calling `flush()` and the `commitTransaction(ignoreError)` method so that it clears the latest error produced by `send(ProducerRecord)`. In order to keep it safe, the change is applied to transactions with no pending record. It mean that, if Therefore, if the user calls `flush()` after a `send`, the error produced by send (if existing) is cleared.
- `commitTransaction(ignoreError)` is defined so that if it is set to true, it MAY clear the last error. The default is false.
- If the user uses `flush()`+ `commitTransaction(true)` then the latest error is ignored and the transaction transits back from the error state. it aaffects only if `flush()` is called beforehand.
Public Interfaces
adding an input parameter to the `send` method so that the user is able to determine not going to the error state by a poison pill record.
The following list categorizes all types of errors that cause a transaction to fail. The category with a beside specifies the cases that the new `send` API is going to cover.
- producer-side errors
- recoverable, such as `RecordTooLargeException`
- irrecoverable, such as `ProducerFencedException`
- recoverable, such as `RecordTooLargeException`
- broker-side errors
Currently, producer-side recoverable errors (the KIP's target category) prevent a record from being added to a batch. They additionally make a transition to an `error` state, which causes the transaction to fail. This KIP provides the possibility to commit a transaction successfully in presence of such errors. Obviously, the problematic records are not added to the batch, but the transition to the `error` state is not done. In other words, with the new `send` API, the transaction does not fail because of a single poison pill record that is not even present in the batch. Obviously, the transaction can still fail due to other types of errors (for example, broker-side errors).
Public Interfaces
If the user 1) is performing a transaction and 2) passes the `TxnSendOption` with the value `IGNORE_SEND_ERRORS` to the `send` method, any poison pill record is excluded from the batch, and the transaction is committed successfully. Note that if the user sets the `TxnSendOption` to `IGNORE_SEND_ERRORS` outside of a transaction, the overloaded `send` method throws an `IllegalStateException`.
| Code Block | ||||
|---|---|---|---|---|
| ||||
/** public Future<RecordMetadata> send(ProducerRecord<K, V> record, *Callback Ifcallback, {@linkTxnSendOption #flush(option) {} is called explicitly before thispublic methodenum andTxnSendOption the{ input config parameter determines ignoring errors, /** * this method clears the last exception produced* byThe theirrecoverable {@link #send(ProducerRecord)} callerrors and transitslead the transaction * out of to the error state. Thereupon, thewhich methodends carriesup outan theunsuccessful samecommit. procedures as delineated in {@link #commitTransaction()}. * / * @param commitOptions The method options NONE, * @throws IllegalStateException if/** no transactional.id has been configured or no transaction has* beenThe started records causing irrecoverable errors are *excluded @throwsfrom ProducerFencedExceptionthe fatal error indicating another producer with the same transactional.idbatch and the transaction is active * @throws org.apache.kafka.common.errors.UnsupportedVersionException fatal error indicating the broker committed successfully. * Note to use this, does not supportonly in transactions (i.e. ifOtherwise its version is lower than 0.11.0.0) * @throws org.apache.kafka.common.errors.AuthorizationException fatal error indicating that the configured * transactional.id is not authorized. See the exception for more details * @throws org.apache.kafka.common.errors.InvalidProducerEpochException if the producer has attempted to produce with an old epoch * to the partition leader. See the exception for more details * @throws KafkaException if the producer has encountered a previous fatal or abortable error, or for any * other unexpected error{@link #send(ProducerRecord, Callback, TxnSendOption)} throws exception. */ IGNORE_SEND_ERRORS } |
We define the method in the interface as well. The `default` implementation helps with backward compatibility.
| Code Block | ||||
|---|---|---|---|---|
| ||||
/** * @throwsSee TimeoutException if the time taken for committing the transaction has surpassed <code>max.block.ms</code>. * @throws InterruptException if the thread is interrupted while blocked {@link KafkaProducer#send(ProducerRecord, Callback, SendTxnOption)} */ default public void commitTransaction(Map<String, ?> commitOptionssend(ProducerRecord record, Callback callback, SendTxnOption option) throws ProducerFencedException {} |
Compatibility, Deprecation, and Migration Plan
Since the default behaviour is preserved, the change has no impact on existing users.
Test Plan
Unit tests for `KafkaProducer` to show that the new feature works with different `send()` errors and exceptions such as RecordTooLargeException.
Rejected Alternatives
...
- Add a producer custom exception handler interface: see KIP-1038 and the discussions.
- Identify poison pill records application-side: It is not efficient and sometimes not even doable. For example for identifying too-large-records the application must be aware of producer configs as well as serialization method, which is not feasible sometimes. More over, checking every single record's size before sending it to catch the bad record is an overhead considering that this check is done by Producer as well.
- Add the feature of clearing errors to `flush()`: `flush()` is not necessarily a transactional method + `AddPartition` is not done successfully in the next `send`.
- Add the feature of clearing errors to `commitTxn()`: `AddPartition` is not done successfully in the next `send`.Just `flush()
- `send()` throws ApiException: may break backward compatibility.
- No need of explicit `flush()` before calling `Just
commitTransaction(commitOptions)`: not safe + `AddPartition` is not done successfully in the next `send` - The user must use ALOS instead of EOS in case they want to drop the poison pill records: ALOS can not guarantee not having duplicates.