spring-cloud / spring-cloud/spring-cloud-function

KafkaNull payload bypasses MessageConverterHelper and is passed unconverted to a typed Consumer, causing ClassCastException

Open
#1,448 0 comments 0 reactions 0 assignees View on GitHub

Nobody has claimed this yet.

Dominant language
Java
Stars
1.1k
Forks
641
Avg merge
11h 2m
Merged PRs (30d)
8

Description

Describe the bug

When a Kafka consumer function is bound to a concrete argument type (e.g. Consumer<MyDomainType>), and the broker delivers a record with a null value — whether from log-compaction tombstone semantics, or simply a producer publishing an empty/null payload — Spring Cloud Stream's Kafka binder wraps it as a Message whose payload is KafkaNull.INSTANCE. SimpleFunctionRegistry.FunctionInvocationWrapper#convertInputIfNecessary detects this and returns the unconverted, raw Message as the function's input argument, instead of attempting conversion or signalling a conversion failure:

https://github.com/spring-cloud/spring-cloud-function/blob/main/spring-cloud-function-context/src/main/java/org/springframework/cloud/function/context/catalog/SimpleFunctionRegistry.java

else if (input instanceof Message) {
    input = this.filterOutHeaders((Message) input);
    if (((Message) input).getPayload().getClass().getName().equals("org.springframework.kafka.support.KafkaNull")) {
        return input;   // bypasses every converter and the MessageConverterHelper hook entirely
    }
    ...

The generated bridge method for the functional interface then tries to cast that Message to the function's declared argument type, throwing a bare ClassCastException from SimpleFunctionRegistry.FunctionInvocationWrapper#invokeConsumer.

Critically, this happens without ever calling any registered MessageConverter or MessageConverterHelper, so the extension point added to resolve #1168 (letting an application signal "this should be a fatal conversion failure, not passed to my function") has no effect here — the KafkaNull branch returns early before that machinery is ever reached. Frameworks and applications built on top of this therefore see a plain ClassCastException with no way to distinguish "the framework couldn't convert this" from "a bug in my own consumer code."

This is closely related to #1168, which covers the case where a registered converter throws during conversion. This issue covers a narrower, still-open gap: the KafkaNull payload never reaches a converter at all, so #1168's fix doesn't help here.

Environment

  • spring-cloud-function-context 4.3.4
  • spring-cloud-stream / spring-cloud-stream-binder-kafka 4.3.3
  • spring-kafka 3.3.16
  • (the shortcut appears structurally unchanged across the versions we checked, so this likely affects a wide range)

To Reproduce

  1. Define a functional Kafka consumer bound to a concrete type, e.g. Consumer<Event<SomeDomainType>>.
  2. Publish a Kafka record with a null value to the bound topic.
  3. Observe: the consumer function is invoked with a raw Message object instead of the declared type, throwing ClassCastException (e.g. class org.springframework.messaging.support.GenericMessage cannot be cast to class com.example.SomeDomainType).

We can provide a minimal standalone Spring Boot + Testcontainers-Kafka reproduction project if useful.

Expected behavior

When the payload is KafkaNull and the target function's declared argument type is not itself Message/KafkaNull-compatible, the framework should route through the same MessageConverterHelper.shouldFailIfCantConvert mechanism used elsewhere — either raising a MessageConversionException or letting the registered helper decide — rather than unconditionally passing the raw Message through. A function genuinely declared to accept Message<?>/KafkaNull (a real null-value-aware consumer) should keep working exactly as it does today; only the mismatched-type case should change.

Suggested fix

In the KafkaNull branch of convertInputIfNecessary, check whether the function's declared parameter type is assignable from Message/KafkaNull before returning the raw input. If it isn't, construct a MessageConversionException describing the failed conversion and call messageConverterHelper.shouldFailIfCantConvert(message, thatException) — passing the actual exception, not null (note: shouldFailIfCantConvert short-circuits to false whenever the throwable argument is null, both in the default implementation and in ours, so the framework needs to construct and pass a real exception here, not signal with a null throwable). If the helper returns true, throw the exception instead of returning the raw message. This reuses the exact mechanism #1168 introduced rather than adding a new one.

Current workaround

We've worked around this downstream, in our own Kafka error-handling layer, by pattern-matching the resulting ClassCastException's message (checking that the class being cast from is GenericMessage) and reclassifying it as a MessageConversionException after the fact. It works, but it's inherently fragile (string-matching an exception message), and only possible because we own that error-handling layer — most Spring Cloud Stream users wouldn't have that option.

Contributor guide

Open the contributing guide

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

Research direction

Start in spring-cloud-function-context/src/main/java/org/springframework/cloud/function/context/catalog/SimpleFunctionRegistry.java, at FunctionInvocationWrapper#convertInputIfNecessary and the KafkaNull early return. Compare that path with the MessageConverterHelper flow introduced for #1168, then reproduce it with a typed Kafka consumer receiving a null record. Done means incompatible KafkaNull payloads produce the intended conversion failure while Message/KafkaNull-compatible consumers retain current behavior.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, kafka, spring
Domain
backend
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Active
Clarity
Clearly specified
Newbie friendliness
55/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.