diff --git a/topic/src/main/java/tech/ydb/topic/description/CodecRegistry.java b/topic/src/main/java/tech/ydb/topic/description/CodecRegistry.java index 513201770..04df98b78 100644 --- a/topic/src/main/java/tech/ydb/topic/description/CodecRegistry.java +++ b/topic/src/main/java/tech/ydb/topic/description/CodecRegistry.java @@ -1,7 +1,7 @@ package tech.ydb.topic.description; -import java.util.HashMap; import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -17,10 +17,11 @@ public class CodecRegistry { private static final Logger logger = LoggerFactory.getLogger(CodecRegistry.class); - final Map customCodecMap; + // registerCodec may be called at any moment, while getCodec is used by the compression and the + // decompression threads of every reader and writer created from the same TopicClient + private final Map customCodecMap = new ConcurrentHashMap<>(); public CodecRegistry() { - customCodecMap = new HashMap<>(); for (Codec codec: StandardCodecs.getAvailableCodecs()) { customCodecMap.put(codec.getId(), codec); } @@ -32,11 +33,12 @@ public CodecRegistry() { * @return previous implementation with associated codec */ public Codec registerCodec(Codec codec) { - assert codec != null; + if (codec == null) { + throw new IllegalArgumentException("Codec must be not null"); + } int codecId = codec.getId(); Codec result = customCodecMap.put(codecId, codec); - if (result != null) { logger.info( "Replace codec which have already associated with this id. CodecId: {} Codec: {}", diff --git a/topic/src/test/java/tech/ydb/topic/impl/CodecRegistryTest.java b/topic/src/test/java/tech/ydb/topic/impl/CodecRegistryTest.java index 35e0baa77..b341c1093 100644 --- a/topic/src/test/java/tech/ydb/topic/impl/CodecRegistryTest.java +++ b/topic/src/test/java/tech/ydb/topic/impl/CodecRegistryTest.java @@ -1,5 +1,9 @@ package tech.ydb.topic.impl; +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; + import org.junit.Assert; import org.junit.Before; import org.junit.Test; @@ -7,10 +11,6 @@ import tech.ydb.topic.description.Codec; import tech.ydb.topic.description.CodecRegistry; -import java.io.IOException; -import java.io.InputStream; -import java.io.OutputStream; - /** * Unit tests for check simple logic for register custom codec * @@ -19,7 +19,7 @@ public class CodecRegistryTest { CodecRegistry registry; - private static final int codecId = 10113; + private static final int CUSTOM_ID = 10113; @Before public void beforeTest() { @@ -34,15 +34,14 @@ public void registerCustomCodecShouldDoubleRegisterCodecAndReturnLastCodec() { registry.registerCodec(codec1); Assert.assertEquals(codec1, registry.registerCodec(codec2)); - Assert.assertEquals(codec2, registry.getCodec(codecId)); - Assert.assertNotEquals(codec1, registry.getCodec(codecId)); + Assert.assertEquals(codec2, registry.getCodec(CUSTOM_ID)); + Assert.assertNotEquals(codec1, registry.getCodec(CUSTOM_ID)); } @Test public void registerCustomCodecShouldNotAcceptNull() { - Assert.assertThrows( - AssertionError.class, - () -> registry.registerCodec(null)); + Exception ex = Assert.assertThrows(IllegalArgumentException.class, () -> registry.registerCodec(null)); + Assert.assertEquals("Codec must be not null", ex.getMessage()); } @Test @@ -65,7 +64,12 @@ static class CodecTopic implements Codec { int codec; public CodecTopic() { - this.codec = codecId; + this.codec = CUSTOM_ID; + } + + @Override + public String toString() { + return "CustomCodec"; } public void setCodecId(int codecId) {