diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java index 2a857c4f25451..e5299cb12b5a6 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/SchemasResourceBase.java @@ -23,6 +23,7 @@ import static org.apache.commons.lang.StringUtils.defaultIfEmpty; import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Charsets; +import java.io.IOException; import java.nio.ByteBuffer; import java.time.Clock; import java.util.List; @@ -143,8 +144,14 @@ public void postSchema(PostSchemaPayload payload, boolean authoritative, AsyncRe } byte[] data; if (SchemaType.KEY_VALUE.name().equals(payload.getType())) { - data = DefaultImplementation - .convertKeyValueDataStringToSchemaInfoSchema(payload.getSchema().getBytes(Charsets.UTF_8)); + try { + data = DefaultImplementation.getDefaultImplementation() + .convertKeyValueDataStringToSchemaInfoSchema(payload.getSchema().getBytes(Charsets.UTF_8)); + } catch (IOException conversionError) { + log.error("[{}] Failed to post schema for topic {}", clientAppId(), topicName, conversionError); + response.resume(Response.serverError().build()); + return; + } } else { data = payload.getSchema().getBytes(Charsets.UTF_8); } @@ -241,16 +248,21 @@ protected String domain() { } private static GetSchemaResponse convertSchemaAndMetadataToGetSchemaResponse(SchemaAndMetadata schemaAndMetadata) { - String schemaData; - if (schemaAndMetadata.schema.getType() == SchemaType.KEY_VALUE) { - schemaData = DefaultImplementation.convertKeyValueSchemaInfoDataToString( - DefaultImplementation.decodeKeyValueSchemaInfo(schemaAndMetadata.schema.toSchemaInfo())); - } else { - schemaData = new String(schemaAndMetadata.schema.getData(), UTF_8); + try { + String schemaData; + if (schemaAndMetadata.schema.getType() == SchemaType.KEY_VALUE) { + schemaData = DefaultImplementation.getDefaultImplementation().convertKeyValueSchemaInfoDataToString( + DefaultImplementation.getDefaultImplementation() + .decodeKeyValueSchemaInfo(schemaAndMetadata.schema.toSchemaInfo())); + } else { + schemaData = new String(schemaAndMetadata.schema.getData(), UTF_8); + } + return GetSchemaResponse.builder().version(getLongSchemaVersion(schemaAndMetadata.version)) + .type(schemaAndMetadata.schema.getType()).timestamp(schemaAndMetadata.schema.getTimestamp()) + .data(schemaData).properties(schemaAndMetadata.schema.getProps()).build(); + } catch (IOException conversionError) { + throw new RuntimeException(conversionError); } - return GetSchemaResponse.builder().version(getLongSchemaVersion(schemaAndMetadata.version)) - .type(schemaAndMetadata.schema.getType()).timestamp(schemaAndMetadata.schema.getTimestamp()) - .data(schemaData).properties(schemaAndMetadata.schema.getProps()).build(); } private static void handleGetSchemaResponse(AsyncResponse response, SchemaAndMetadata schema, Throwable error) { diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessageIdSerializationTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessageIdSerializationTest.java index 01661ce84eea2..295c7803372f9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessageIdSerializationTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/MessageIdSerializationTest.java @@ -19,7 +19,7 @@ package org.apache.pulsar.broker.service; import static org.testng.Assert.assertEquals; - +import java.io.IOException; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.impl.MessageIdImpl; import org.testng.annotations.Test; @@ -46,7 +46,7 @@ public void testProtobufSerializationNull() throws Exception { MessageId.fromByteArray(null); } - @Test(expectedExceptions = RuntimeException.class) + @Test(expectedExceptions = IOException.class) void testProtobufSerializationEmpty() throws Exception { MessageId.fromByteArray(new byte[0]); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TransactionEndToEndTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TransactionEndToEndTest.java index f971b5a770cca..d7cb6c9cc3220 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TransactionEndToEndTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TransactionEndToEndTest.java @@ -66,7 +66,6 @@ import org.apache.pulsar.common.naming.NamespaceName; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.ClusterData; -import org.apache.pulsar.common.policies.data.ClusterDataImpl; import org.apache.pulsar.common.policies.data.PersistentTopicInternalStats; import org.apache.pulsar.common.policies.data.TenantInfoImpl; import org.apache.pulsar.common.util.collections.ConcurrentOpenHashMap; @@ -548,7 +547,7 @@ private void txnCumulativeAckTest(boolean batchEnable, int maxBatchSize, Subscri } try { - consumer.acknowledgeCumulativeAsync(DefaultImplementation + consumer.acknowledgeCumulativeAsync(DefaultImplementation.getDefaultImplementation() .newMessageId(((MessageIdImpl) message.getMessageId()).getLedgerId(), ((MessageIdImpl) message.getMessageId()).getEntryId() - 1, -1), abortTxn).get(); diff --git a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/SchemasImpl.java b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/SchemasImpl.java index 4408ae21b50f0..a072acd8b73a4 100644 --- a/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/SchemasImpl.java +++ b/pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/SchemasImpl.java @@ -19,6 +19,7 @@ package org.apache.pulsar.client.admin.internal; import static java.nio.charset.StandardCharsets.UTF_8; +import java.io.IOException; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; @@ -438,8 +439,12 @@ static SchemaInfo convertGetSchemaResponseToSchemaInfo(TopicName tn, byte[] schema; if (response.getType() == SchemaType.KEY_VALUE) { - schema = DefaultImplementation.convertKeyValueDataStringToSchemaInfoSchema( - response.getData().getBytes(UTF_8)); + try { + schema = DefaultImplementation.getDefaultImplementation().convertKeyValueDataStringToSchemaInfoSchema( + response.getData().getBytes(UTF_8)); + } catch (IOException conversionError) { + throw new RuntimeException(conversionError); + } } else { schema = response.getData().getBytes(UTF_8); } @@ -466,28 +471,31 @@ static SchemaInfoWithVersion convertGetSchemaResponseToSchemaInfoWithVersion(Top // the util function exists for backward compatibility concern - static String convertSchemaDataToStringLegacy(SchemaInfo schemaInfo) { + static String convertSchemaDataToStringLegacy(SchemaInfo schemaInfo) throws IOException { byte[] schemaData = schemaInfo.getSchema(); if (null == schemaInfo.getSchema()) { return ""; } if (schemaInfo.getType() == SchemaType.KEY_VALUE) { - return DefaultImplementation.convertKeyValueSchemaInfoDataToString( - DefaultImplementation.decodeKeyValueSchemaInfo(schemaInfo)); + return DefaultImplementation.getDefaultImplementation().convertKeyValueSchemaInfoDataToString( + DefaultImplementation.getDefaultImplementation().decodeKeyValueSchemaInfo(schemaInfo)); } return new String(schemaData, UTF_8); } static PostSchemaPayload convertSchemaInfoToPostSchemaPayload(SchemaInfo schemaInfo) { - - PostSchemaPayload payload = new PostSchemaPayload(); - payload.setType(schemaInfo.getType().name()); - payload.setProperties(schemaInfo.getProperties()); - // for backward compatibility concern, we convert `bytes` to `string` - // we can consider fixing it in a new version of rest endpoint - payload.setSchema(convertSchemaDataToStringLegacy(schemaInfo)); - return payload; + try { + PostSchemaPayload payload = new PostSchemaPayload(); + payload.setType(schemaInfo.getType().name()); + payload.setProperties(schemaInfo.getProperties()); + // for backward compatibility concern, we convert `bytes` to `string` + // we can consider fixing it in a new version of rest endpoint + payload.setSchema(convertSchemaDataToStringLegacy(schemaInfo)); + return payload; + } catch (IOException conversionError) { + throw new RuntimeException(conversionError); + } } } diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/AuthenticationFactory.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/AuthenticationFactory.java index a244c3d2c8e4c..4de12b2eff2f0 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/AuthenticationFactory.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/AuthenticationFactory.java @@ -41,7 +41,7 @@ public final class AuthenticationFactory { * @return the Authentication object initialized with the token credentials */ public static Authentication token(String token) { - return DefaultImplementation.newAuthenticationToken(token); + return DefaultImplementation.getDefaultImplementation().newAuthenticationToken(token); } /** @@ -52,7 +52,7 @@ public static Authentication token(String token) { * @return the Authentication object initialized with the token credentials */ public static Authentication token(Supplier tokenSupplier) { - return DefaultImplementation.newAuthenticationToken(tokenSupplier); + return DefaultImplementation.getDefaultImplementation().newAuthenticationToken(tokenSupplier); } // CHECKSTYLE.OFF: MethodName @@ -67,7 +67,7 @@ public static Authentication token(Supplier tokenSupplier) { * @return the Authentication object initialized with the TLS credentials */ public static Authentication TLS(String certFilePath, String keyFilePath) { - return DefaultImplementation.newAuthenticationTLS(certFilePath, keyFilePath); + return DefaultImplementation.getDefaultImplementation().newAuthenticationTLS(certFilePath, keyFilePath); } // CHECKSTYLE.ON: MethodName @@ -86,7 +86,7 @@ public static Authentication TLS(String certFilePath, String keyFilePath) { public static Authentication create(String authPluginClassName, String authParamsString) throws UnsupportedAuthenticationException { try { - return DefaultImplementation.createAuthentication(authPluginClassName, authParamsString); + return DefaultImplementation.getDefaultImplementation().createAuthentication(authPluginClassName, authParamsString); } catch (Throwable t) { throw new UnsupportedAuthenticationException(t); } @@ -103,7 +103,7 @@ public static Authentication create(String authPluginClassName, String authParam public static Authentication create(String authPluginClassName, Map authParams) throws UnsupportedAuthenticationException { try { - return DefaultImplementation.createAuthentication(authPluginClassName, authParams); + return DefaultImplementation.getDefaultImplementation().createAuthentication(authPluginClassName, authParams); } catch (Throwable t) { throw new UnsupportedAuthenticationException(t); } diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/BatcherBuilder.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/BatcherBuilder.java index d659a952e8fcb..34d8375fbcead 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/BatcherBuilder.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/BatcherBuilder.java @@ -19,6 +19,7 @@ package org.apache.pulsar.client.api; import java.io.Serializable; + import org.apache.pulsar.client.internal.DefaultImplementation; import org.apache.pulsar.common.classification.InterfaceAudience; import org.apache.pulsar.common.classification.InterfaceStability; @@ -39,7 +40,7 @@ public interface BatcherBuilder extends Serializable { *

batched into single batch message: * [(k1, v1), (k2, v1), (k3, v1), (k1, v2), (k2, v2), (k3, v2), (k1, v3), (k2, v3), (k3, v3)] */ - BatcherBuilder DEFAULT = DefaultImplementation.newDefaultBatcherBuilder(); + BatcherBuilder DEFAULT = DefaultImplementation.getDefaultImplementation().newDefaultBatcherBuilder(); /** * Key based batch message container. @@ -50,7 +51,7 @@ public interface BatcherBuilder extends Serializable { *

batched into multiple batch messages: * [(k1, v1), (k1, v2), (k1, v3)], [(k2, v1), (k2, v2), (k2, v3)], [(k3, v1), (k3, v2), (k3, v3)] */ - BatcherBuilder KEY_BASED = DefaultImplementation.newKeyBasedBatcherBuilder(); + BatcherBuilder KEY_BASED = DefaultImplementation.getDefaultImplementation().newKeyBasedBatcherBuilder(); /** * Build a new batch message container. diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/MessageId.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/MessageId.java index 7b275c448d3df..bf8defefe0c61 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/MessageId.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/MessageId.java @@ -20,6 +20,7 @@ import java.io.IOException; import java.io.Serializable; + import org.apache.pulsar.client.internal.DefaultImplementation; import org.apache.pulsar.common.classification.InterfaceAudience; import org.apache.pulsar.common.classification.InterfaceStability; @@ -54,7 +55,7 @@ public interface MessageId extends Comparable, Serializable { * @throws IOException if the de-serialization fails */ static MessageId fromByteArray(byte[] data) throws IOException { - return DefaultImplementation.newMessageIdFromByteArray(data); + return DefaultImplementation.getDefaultImplementation().newMessageIdFromByteArray(data); } /** @@ -70,7 +71,7 @@ static MessageId fromByteArray(byte[] data) throws IOException { * @throws IOException if the de-serialization fails */ static MessageId fromByteArrayWithTopic(byte[] data, String topicName) throws IOException { - return DefaultImplementation.newMessageIdFromByteArrayWithTopic(data, topicName); + return DefaultImplementation.getDefaultImplementation().newMessageIdFromByteArrayWithTopic(data, topicName); } // CHECKSTYLE.OFF: ConstantName @@ -78,12 +79,12 @@ static MessageId fromByteArrayWithTopic(byte[] data, String topicName) throws IO /** * MessageId that represents the oldest message available in the topic. */ - MessageId earliest = DefaultImplementation.newMessageId(-1, -1, -1); + MessageId earliest = DefaultImplementation.getDefaultImplementation().newMessageId(-1, -1, -1); /** * MessageId that represents the next message published in the topic. */ - MessageId latest = DefaultImplementation.newMessageId(Long.MAX_VALUE, Long.MAX_VALUE, -1); + MessageId latest = DefaultImplementation.getDefaultImplementation().newMessageId(Long.MAX_VALUE, Long.MAX_VALUE, -1); // CHECKSTYLE.ON: ConstantName } diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/PulsarClient.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/PulsarClient.java index 5cf53090455b4..5f1a18b0e0235 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/PulsarClient.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/PulsarClient.java @@ -52,7 +52,7 @@ public interface PulsarClient extends Closeable { * @since 2.0.0 */ static ClientBuilder builder() { - return DefaultImplementation.newClientBuilder(); + return DefaultImplementation.getDefaultImplementation().newClientBuilder(); } /** diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/Schema.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/Schema.java index 6cf9ccdc794be..6803187521adf 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/Schema.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/Schema.java @@ -18,7 +18,7 @@ */ package org.apache.pulsar.client.api; -import static org.apache.pulsar.client.internal.DefaultImplementation.getBytes; +import static org.apache.pulsar.client.internal.PulsarClientImplementationBinding.getBytes; import java.nio.ByteBuffer; import java.sql.Time; import java.sql.Timestamp; @@ -174,7 +174,7 @@ default void configureSchemaInfo(String topic, String componentName, /** * Schema that doesn't perform any encoding on the message payloads. Accepts a byte array and it passes it through. */ - Schema BYTES = DefaultImplementation.newBytesSchema(); + Schema BYTES = DefaultImplementation.getDefaultImplementation().newBytesSchema(); /** * Return the native schema that is wrapped by Pulsar API. @@ -188,79 +188,79 @@ default Optional getNativeSchema() { /** * ByteBuffer Schema. */ - Schema BYTEBUFFER = DefaultImplementation.newByteBufferSchema(); + Schema BYTEBUFFER = DefaultImplementation.getDefaultImplementation().newByteBufferSchema(); /** * Schema that can be used to encode/decode messages whose values are String. The payload is encoded with UTF-8. */ - Schema STRING = DefaultImplementation.newStringSchema(); + Schema STRING = DefaultImplementation.getDefaultImplementation().newStringSchema(); /** * INT8 Schema. */ - Schema INT8 = DefaultImplementation.newByteSchema(); + Schema INT8 = DefaultImplementation.getDefaultImplementation().newByteSchema(); /** * INT16 Schema. */ - Schema INT16 = DefaultImplementation.newShortSchema(); + Schema INT16 = DefaultImplementation.getDefaultImplementation().newShortSchema(); /** * INT32 Schema. */ - Schema INT32 = DefaultImplementation.newIntSchema(); + Schema INT32 = DefaultImplementation.getDefaultImplementation().newIntSchema(); /** * INT64 Schema. */ - Schema INT64 = DefaultImplementation.newLongSchema(); + Schema INT64 = DefaultImplementation.getDefaultImplementation().newLongSchema(); /** * Boolean Schema. */ - Schema BOOL = DefaultImplementation.newBooleanSchema(); + Schema BOOL = DefaultImplementation.getDefaultImplementation().newBooleanSchema(); /** * Float Schema. */ - Schema FLOAT = DefaultImplementation.newFloatSchema(); + Schema FLOAT = DefaultImplementation.getDefaultImplementation().newFloatSchema(); /** * Double Schema. */ - Schema DOUBLE = DefaultImplementation.newDoubleSchema(); + Schema DOUBLE = DefaultImplementation.getDefaultImplementation().newDoubleSchema(); /** * Date Schema. */ - Schema DATE = DefaultImplementation.newDateSchema(); + Schema DATE = DefaultImplementation.getDefaultImplementation().newDateSchema(); /** * Time Schema. */ - Schema