Azure / Azure/azure-sdk-for-java

[Cosmos Kafka Connector] ITEM_PATCH does not skip filter-predicate 412s (Spark connector does, after #49700)

Open
#49,944 1 comment 1 reaction 1 assignee Claimed by @tvaron3 View on GitHub
bug Client Cosmos Service Attention
Dominant language
Java
Stars
2.6k
Forks
2.2k
Avg merge
2d 8h
Merged PRs (30d)
178

Description

### Summary

The Cosmos DB **Spark** connector was changed in #49700 so that an HTTP `412 Precondition Failed` — returned by the service for a document excluded by a conditional patch filter (`spark.cosmos.write.patch.filter`) — is treated as a successful no-op skip rather than failing the write.

The **Kafka** sink connector supports the equivalent config (`azure.cosmos.sink.write.patch.filter`, see `sdk/cosmos/azure-cosmos-kafka-connect/docs/configuration-reference.md:67`) and has the same `shouldIgnore`-by-write-strategy structure, but `ITEM_PATCH` still falls through to `default: return false` and therefore fails on a filter-predicate 412.

### Where

`sdk/cosmos/azure-cosmos-kafka-connect/src/main/java/com/azure/cosmos/kafka/connect/implementation/sink/CosmosBulkWriter.java` (~line 320):

```java
private boolean shouldIgnore(BulkOperationFailedException failedException) {
switch (this.writeConfig.getItemWriteStrategy()) {
case ITEM_APPEND:
return KafkaCosmosExceptionsHelper.isResourceExistsException(failedException);
case ITEM_DELETE:
return KafkaCosmosExceptionsHelper.isNotFoundException(failedException);
case ITEM_DELETE_IF_NOT_MODIFIED:
return KafkaCosmosExceptionsHelper.isNotFoundException(failedException)
|| KafkaCosmosExceptionsHelper.isPreconditionFailedException(failedException);
case ITEM_OVERWRITE_IF_NOT_MODIFIED:
return KafkaCosmosExceptionsHelper.isResourceExistsException(failedException)
|| KafkaCosmosExceptionsHelper.isNotFoundException(failedException)
|| KafkaCosmosExceptionsHelper.isPreconditionFailedException(failedException);
default: // <-- ITEM_PATCH lands here
return false;
}
}
```

`CosmosPointWriter.java` (`case ITEM_PATCH:` ~line 68) has the same gap on the point-write path.

### Why it matters

The customer scenario that motivated #49700 (see #49594) is a **Kafka -> Spark -> Cosmos** pipeline doing idempotent server-side `increment` patches guarded by a filter predicate such as `NOT IS_DEFINED(last_batch_id) OR last_batch_id < `. On replay the predicate legitimately excludes already-applied documents, producing 412s.

A customer who moves that same logic to the Kafka sink connector directly — a very natural step, since it removes Spark from the pipeline — will hit exactly the bug that #49700 just fixed for Spark. The two connectors share the customer scenario, so the divergence is likely to be discovered by a customer rather than by us.

### Proposed fix

Mirror the Spark behavior: treat a 412 as an ignorable no-op for `ITEM_PATCH` when a non-empty patch filter predicate is configured, in both `CosmosBulkWriter.shouldIgnore` and the `ITEM_PATCH` path of `CosmosPointWriter`.

Rationale (same as #49700): a patch filter predicate's contract is "only modify the document when the condition holds, otherwise leave it alone", so a predicate miss is an expected outcome rather than an error. This is also consistent with the connector's own existing unconditional 412-ignore for `ITEM_DELETE_IF_NOT_MODIFIED` and `ITEM_OVERWRITE_IF_NOT_MODIFIED`.

Note that etag / `If-Match` is not wired for the patch path in either connector today, so the filter predicate is currently the only source of a 412 on patch. If etag support is added later, the two 412 causes will need to be distinguished — worth carrying the same inline comment the Spark connector now has.

### Notes

- Found during deep review of #49700.
- Consider also mirroring the skip observability added in #49700 (a `recordsSkipped` metric plus a summary log at flush), so a fully-skipped run is distinguishable from a fully-successful one.

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.