Skip to content

Commit e2a8cb3

Browse files
authored
GH-4442: Fix value serializer mappings for different classloaders
Fixes: #4442 When the type mapper is initialized under one classloader and messages are produced under another (e.g. Spring DevTools restart classloader), the reverse lookup map keyed on `Class<?>` identity misses the entry because two `Class<?>` objects representing the same class loaded by different classloaders are not equal in a `HashMap` lookup. The fallback writes the FQCN header instead of the configured alias, silently breaking consumers that depend on the alias. Fix by changing the reverse map key from `Class<?>` to `String` (using `clazz.getName()`), which is classloader-agnostic. The same change is applied to the deprecated `AbstractJavaTypeMapper` for Jackson 2. Signed-off-by: Soby Chacko <soby.chacko@broadcom.com> **Auto-cherry-pick to `4.0.x`**
1 parent 9a85194 commit e2a8cb3

3 files changed

Lines changed: 32 additions & 8 deletions

File tree

spring-kafka/src/main/java/org/springframework/kafka/support/mapping/AbstractJavaTypeMapper.java

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -79,7 +79,7 @@ public abstract class AbstractJavaTypeMapper implements BeanClassLoaderAware {
7979

8080
private final Map<String, Class<?>> idClassMapping = new ConcurrentHashMap<String, Class<?>>();
8181

82-
private final Map<Class<?>, byte[]> classIdMapping = new ConcurrentHashMap<Class<?>, byte[]>();
82+
private final Map<String, byte[]> classIdMapping = new ConcurrentHashMap<String, byte[]>();
8383

8484
private String classIdFieldName = DEFAULT_CLASSID_FIELD_NAME;
8585

@@ -143,8 +143,9 @@ public void setBeanClassLoader(ClassLoader classLoader) {
143143
}
144144

145145
protected void addHeader(Headers headers, String headerName, Class<?> clazz) {
146-
if (this.classIdMapping.containsKey(clazz)) {
147-
headers.add(new RecordHeader(headerName, this.classIdMapping.get(clazz)));
146+
byte[] alias = this.classIdMapping.get(clazz.getName());
147+
if (alias != null) {
148+
headers.add(new RecordHeader(headerName, alias));
148149
}
149150
else {
150151
headers.add(new RecordHeader(headerName, clazz.getName().getBytes(StandardCharsets.UTF_8)));
@@ -177,7 +178,7 @@ private void createReverseMap() {
177178
for (Map.Entry<String, Class<?>> entry : this.idClassMapping.entrySet()) {
178179
String id = entry.getKey();
179180
Class<?> clazz = entry.getValue();
180-
this.classIdMapping.put(clazz, id.getBytes(StandardCharsets.UTF_8));
181+
this.classIdMapping.put(clazz.getName(), id.getBytes(StandardCharsets.UTF_8));
181182
}
182183
}
183184

spring-kafka/src/main/java/org/springframework/kafka/support/mapping/DefaultJacksonJavaTypeMapper.java

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -89,7 +89,7 @@ public class DefaultJacksonJavaTypeMapper implements JacksonJavaTypeMapper, Bean
8989

9090
private final Map<String, Class<?>> idClassMapping = new ConcurrentHashMap<String, Class<?>>();
9191

92-
private final Map<Class<?>, byte[]> classIdMapping = new ConcurrentHashMap<Class<?>, byte[]>();
92+
private final Map<String, byte[]> classIdMapping = new ConcurrentHashMap<String, byte[]>();
9393

9494
private String classIdFieldName = DEFAULT_CLASSID_FIELD_NAME;
9595

@@ -153,8 +153,9 @@ public void setBeanClassLoader(ClassLoader classLoader) {
153153
}
154154

155155
protected void addHeader(Headers headers, String headerName, Class<?> clazz) {
156-
if (this.classIdMapping.containsKey(clazz)) {
157-
headers.add(new RecordHeader(headerName, this.classIdMapping.get(clazz)));
156+
byte[] alias = this.classIdMapping.get(clazz.getName());
157+
if (alias != null) {
158+
headers.add(new RecordHeader(headerName, alias));
158159
}
159160
else {
160161
headers.add(new RecordHeader(headerName, clazz.getName().getBytes(StandardCharsets.UTF_8)));
@@ -187,7 +188,7 @@ private void createReverseMap() {
187188
for (Map.Entry<String, Class<?>> entry : this.idClassMapping.entrySet()) {
188189
String id = entry.getKey();
189190
Class<?> clazz = entry.getValue();
190-
this.classIdMapping.put(clazz, id.getBytes(StandardCharsets.UTF_8));
191+
this.classIdMapping.put(clazz.getName(), id.getBytes(StandardCharsets.UTF_8));
191192
}
192193
}
193194

spring-kafka/src/test/java/org/springframework/kafka/support/serializer/JsonSerializationTests.java

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,9 @@
1717
package org.springframework.kafka.support.serializer;
1818

1919
import java.io.IOException;
20+
import java.net.URL;
21+
import java.net.URLClassLoader;
22+
import java.nio.charset.StandardCharsets;
2023
import java.util.Arrays;
2124
import java.util.Collections;
2225
import java.util.HashMap;
@@ -60,6 +63,7 @@
6063
* @author Torsten Schleede
6164
* @author Gary Russell
6265
* @author Ivan Ponomarev
66+
* @author Soby Chacko
6367
*/
6468
public class JsonSerializationTests {
6569

@@ -434,6 +438,24 @@ void configRejectedIgnoredAfterPropertiesSet() {
434438
assertThatIllegalStateException().isThrownBy(() -> ser.configure(configs2, false));
435439
}
436440

441+
@Test
442+
void typeMappingHonoredWhenClassLoadedByDifferentClassLoader() throws Exception {
443+
DefaultJacksonJavaTypeMapper mapper = new DefaultJacksonJavaTypeMapper();
444+
mapper.setIdClassMapping(Map.of("my-alias", Foo.class));
445+
446+
URL codeSourceUrl = Foo.class.getProtectionDomain().getCodeSource().getLocation();
447+
try (URLClassLoader isolatedLoader = new URLClassLoader(new URL[] { codeSourceUrl }, null)) {
448+
Class<?> reloadedClass = isolatedLoader.loadClass(Foo.class.getName());
449+
assertThat(reloadedClass).isNotSameAs(Foo.class);
450+
451+
Headers headers = new RecordHeaders();
452+
mapper.fromClass(reloadedClass, headers);
453+
454+
String typeHeader = new String(headers.lastHeader("__TypeId__").value(), StandardCharsets.UTF_8);
455+
assertThat(typeHeader).isEqualTo("my-alias");
456+
}
457+
}
458+
437459
public static JavaType fooBarJavaType(byte[] data, Headers headers) {
438460
if (data[0] == '{' && data[1] == 'f') {
439461
return TypeFactory.createDefaultInstance().constructType(Foo.class);

0 commit comments

Comments
 (0)