Github user michaelandrepearce commented on a diff in the pull request:

    https://github.com/apache/activemq-artemis/pull/1607#discussion_r146868275
  
    --- Diff: 
integration/activemq-kafka/activemq-kafka-protocols/activemq-kafka-amqp-protocol/src/main/java/org/apache/activemq/artemis/integration/kafka/protocol/amqp/AmqpMessageSerializer.java
 ---
    @@ -0,0 +1,52 @@
    +/*
    + * Licensed to the Apache Software Foundation (ASF) under one or more
    + * contributor license agreements. See the NOTICE file distributed with
    + * this work for additional information regarding copyright ownership.
    + * The ASF licenses this file to You under the Apache License, Version 2.0
    + * (the "License"); you may not use this file except in compliance with
    + * the License. You may obtain a copy of the License at
    + *
    + *     http://www.apache.org/licenses/LICENSE-2.0
    + *
    + * Unless required by applicable law or agreed to in writing, software
    + * distributed under the License is distributed on an "AS IS" BASIS,
    + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
    + * See the License for the specific language governing permissions and
    + * limitations under the License.
    + */
    +package org.apache.activemq.artemis.integration.kafka.protocol.amqp;
    +
    +import java.util.Map;
    +
    +import org.apache.activemq.artemis.api.core.Message;
    +import 
org.apache.activemq.artemis.integration.kafka.protocol.amqp.proton.ProtonMessageSerializer;
    +import org.apache.activemq.artemis.protocol.amqp.broker.AMQPMessage;
    +import 
org.apache.activemq.artemis.protocol.amqp.converter.CoreAmqpConverter;
    +import org.apache.kafka.common.errors.SerializationException;
    +import org.apache.kafka.common.serialization.Serializer;
    +
    +public class AmqpMessageSerializer implements Serializer<Message> {
    +
    +   ProtonMessageSerializer protonMessageSerializer = new 
ProtonMessageSerializer();
    +
    +   @Override
    +   public byte[] serialize(String topic, Message message) {
    +      if (message == null) return null;
    +      try {
    +         AMQPMessage amqpMessage = CoreAmqpConverter.checkAMQP(message);
    +         return protonMessageSerializer.serialize(topic, 
amqpMessage.getProtonMessage());
    --- End diff --
    
    noted.


---

Reply via email to