From 57e6ea96b800040daffcb1a297aa8d425894ebfe Mon Sep 17 00:00:00 2001 From: Andrew Frieze Date: Wed, 15 Nov 2017 14:07:05 -0700 Subject: [PATCH 1/2] Makes source the sole method of obtaining a consumer and adds reader schemas to kafka consumer --- connectors/activemq/build.sbt | 4 +- .../connectors/ActiveMQConnector.java | 8 +- .../test/ActiveMQForkliftConnector.java | 10 -- connectors/kafka/build.sbt | 4 +- .../forklift/connectors/KafkaConnector.java | 100 +++++++++++++----- .../message/ForkliftAvroMessageUtils.java | 30 ++++++ .../producers/KafkaForkliftProducer.java | 40 +++---- .../ForkliftKafkaAvroDeserializer.java | 38 +++++++ .../integration/AvroMessageTests.java | 81 ++++++++++++-- .../server/TestServiceManager.java | 18 ++++ .../test/resources/schemas/AvroMessage.avsc | 5 +- .../resources/schemas/ForkliftMessage.avsc | 1 + core/build.sbt | 2 +- .../connectors/ForkliftConnectorI.java | 4 +- .../java/forklift/source/SchemaResolver.java | 5 + .../source/decorators/GroupedTopic.java | 4 + .../forklift/source/decorators/Queue.java | 5 +- .../forklift/source/decorators/Topic.java | 17 +++ 18 files changed, 290 insertions(+), 86 deletions(-) create mode 100644 connectors/kafka/src/main/java/forklift/message/ForkliftAvroMessageUtils.java create mode 100644 connectors/kafka/src/main/java/forklift/serializers/ForkliftKafkaAvroDeserializer.java create mode 100644 connectors/kafka/src/test/resources/schemas/ForkliftMessage.avsc create mode 100644 core/src/main/java/forklift/source/SchemaResolver.java diff --git a/connectors/activemq/build.sbt b/connectors/activemq/build.sbt index 2b01798..68205ea 100644 --- a/connectors/activemq/build.sbt +++ b/connectors/activemq/build.sbt @@ -2,7 +2,7 @@ organization := "com.github.dcshock" name := "forklift-activemq" -version := "2.0" +version := "2.1" javacOptions ++= Seq("-source", "1.8") @@ -22,7 +22,7 @@ resolvers ++= Seq( ) libraryDependencies ++= Seq( - "com.github.dcshock" % "forklift" % "2.0", + "com.github.dcshock" % "forklift" % "2.2", "org.apache.activemq" % "activemq-client" % "5.14.0", "org.apache.activemq" % "activemq-broker" % "5.14.0", "com.fasterxml.jackson.core" % "jackson-databind" % "2.7.3", diff --git a/connectors/activemq/src/main/java/forklift/connectors/ActiveMQConnector.java b/connectors/activemq/src/main/java/forklift/connectors/ActiveMQConnector.java index 6a629b6..4cb40d6 100644 --- a/connectors/activemq/src/main/java/forklift/connectors/ActiveMQConnector.java +++ b/connectors/activemq/src/main/java/forklift/connectors/ActiveMQConnector.java @@ -86,8 +86,7 @@ public synchronized Session getSession() } } - @Override - public ForkliftConsumerI getQueue(String name) + private ForkliftConsumerI getQueue(String name) throws ConnectorException { final Session s = getSession(); try { @@ -97,8 +96,7 @@ public ForkliftConsumerI getQueue(String name) } } - @Override - public ForkliftConsumerI getTopic(String name) + private ForkliftConsumerI getTopic(String name) throws ConnectorException { final Session s = getSession(); try { @@ -111,7 +109,7 @@ public ForkliftConsumerI getTopic(String name) @Override public ForkliftConsumerI getConsumerForSource(SourceI source) throws ConnectorException { return source - .apply(QueueSource.class, queue -> getQueue(queue.getName())) + .apply(QueueSource.class, queue -> getQueue(queue.getName())) .apply(TopicSource.class, topic -> getTopic(topic.getName())) .apply(GroupedTopicSource.class, topic -> getGroupedTopic(topic)) .apply(RoleInputSource.class, roleSource -> { diff --git a/connectors/activemq/src/test/java/forklift/activemq/test/ActiveMQForkliftConnector.java b/connectors/activemq/src/test/java/forklift/activemq/test/ActiveMQForkliftConnector.java index 72fab75..9f0de98 100644 --- a/connectors/activemq/src/test/java/forklift/activemq/test/ActiveMQForkliftConnector.java +++ b/connectors/activemq/src/test/java/forklift/activemq/test/ActiveMQForkliftConnector.java @@ -34,16 +34,6 @@ public Connection getConnection() throws ConnectorException { return TestServiceManager.getConnector().getConnection(); } - @Override - public ForkliftConsumerI getQueue(String name) throws ConnectorException { - return TestServiceManager.getConnector().getQueue(name); - } - - @Override - public ForkliftConsumerI getTopic(String name) throws ConnectorException { - return TestServiceManager.getConnector().getTopic(name); - } - @Override public ForkliftConsumerI getConsumerForSource(SourceI source) throws ConnectorException { return TestServiceManager.getConnector().getConsumerForSource(source); diff --git a/connectors/kafka/build.sbt b/connectors/kafka/build.sbt index 12017ca..8ae8fce 100644 --- a/connectors/kafka/build.sbt +++ b/connectors/kafka/build.sbt @@ -2,7 +2,7 @@ organization := "com.github.dcshock" name := "forklift-kafka" -version := "2.0" +version := "2.1" //required for some test dependencies scalaVersion := "2.11.7" @@ -26,7 +26,7 @@ resolvers ++= Seq( ) libraryDependencies ++= Seq( - "com.github.dcshock" % "forklift" % "2.0" , + "com.github.dcshock" % "forklift" % "2.2" , "com.fasterxml.jackson.core" % "jackson-databind" % "2.7.3", "com.fasterxml.jackson.datatype" % "jackson-datatype-jsr310" % "2.7.3", "org.apache.kafka" % "kafka-clients" % "0.10.1.1-cp1" exclude("org.slf4j","slf4j-log4j12"), diff --git a/connectors/kafka/src/main/java/forklift/connectors/KafkaConnector.java b/connectors/kafka/src/main/java/forklift/connectors/KafkaConnector.java index 8044d32..112a437 100644 --- a/connectors/kafka/src/main/java/forklift/connectors/KafkaConnector.java +++ b/connectors/kafka/src/main/java/forklift/connectors/KafkaConnector.java @@ -1,24 +1,34 @@ package forklift.connectors; +import com.fasterxml.jackson.databind.DeserializationFeature; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ArrayNode; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; import forklift.consumer.ForkliftConsumerI; import forklift.consumer.KafkaTopicConsumer; import forklift.consumer.wrapper.RoleInputConsumerWrapper; import forklift.controller.KafkaController; +import forklift.decorators.Message; +import forklift.message.ForkliftAvroMessageUtils; import forklift.message.MessageStream; import forklift.producers.ForkliftProducerI; import forklift.producers.KafkaForkliftProducer; +import forklift.serializers.ForkliftKafkaAvroDeserializer; import forklift.source.ActionSource; import forklift.source.LogicalSource; import forklift.source.SourceI; import forklift.source.sources.GroupedTopicSource; -import forklift.source.sources.RoleInputSource; import forklift.source.sources.QueueSource; +import forklift.source.sources.RoleInputSource; import forklift.source.sources.TopicSource; - import io.confluent.kafka.serializers.KafkaAvroDeserializer; import io.confluent.kafka.serializers.KafkaAvroDeserializerConfig; import io.confluent.kafka.serializers.KafkaAvroSerializer; import io.confluent.kafka.serializers.KafkaAvroSerializerConfig; +import org.apache.avro.Schema; +import org.apache.avro.specific.SpecificRecord; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.producer.KafkaProducer; @@ -26,9 +36,14 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.io.IOException; +import java.lang.reflect.Field; +import java.lang.reflect.InvocationTargetException; import java.util.HashMap; +import java.util.HashSet; import java.util.Map; import java.util.Properties; +import java.util.Set; import java.util.concurrent.TimeUnit; /** @@ -101,20 +116,43 @@ private KafkaProducer createKafkaProducer() { return new KafkaProducer(producerProperties); } - private KafkaController createController(String topicName) { + private KafkaController createController(GroupedTopicSource source) { + + ForkliftKafkaAvroDeserializer deserializer = new ForkliftKafkaAvroDeserializer(); + Schema readerSchema = null; + try { + readerSchema = produceReaderSchema(source); + } catch (Exception e) { + log.error("Unable to generate reader schema, falling back to writer schema. Defaults may be lost"); + } + Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaHosts); props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, io.confluent.kafka.serializers.KafkaAvroDeserializer.class); - props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, io.confluent.kafka.serializers.KafkaAvroDeserializer.class); + props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ForkliftKafkaAvroDeserializer.class); + if (readerSchema != null) { + props.put("forklift.avro.reader.schema", readerSchema); + } props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 200); props.put(KafkaAvroDeserializerConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistries); props.put(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, false); final KafkaConsumer kafkaConsumer = new KafkaConsumer(props); - return new KafkaController(kafkaConsumer, new MessageStream(), topicName); + return new KafkaController(kafkaConsumer, new MessageStream(), source.getName()); + } + + private Schema produceReaderSchema(SourceI source) throws NoSuchMethodException, InvocationTargetException, IllegalAccessException, IOException { + Set fields = new HashSet<>(); + for (Field field : source.getContextClass().getDeclaredFields()) { + if (field.isAnnotationPresent(Message.class) && SpecificRecord.class.isAssignableFrom(field.getType())) { + Schema schema = (Schema) field.getType().getMethod("getClassSchema").invoke(null); + return ForkliftAvroMessageUtils.addForkliftPropertiesToSchema(schema); + } + } + return null; } @Override @@ -138,21 +176,30 @@ public synchronized void stop() throws ConnectorException { @Override public ForkliftConsumerI getConsumerForSource(SourceI source) throws ConnectorException { return source - .apply(QueueSource.class, queue -> getQueue(queue.getName())) - .apply(TopicSource.class, topic -> getTopic(topic.getName())) - .apply(GroupedTopicSource.class, topic -> getGroupedTopic(topic)) - .apply(RoleInputSource.class, roleSource -> { - final ForkliftConsumerI rawConsumer = getConsumerForSource(roleSource.getActionSource(this)); - return new RoleInputConsumerWrapper(rawConsumer); - }) - .elseUnsupportedError(); + .apply(QueueSource.class, queue -> { + GroupedTopicSource groupedTopicSource = new GroupedTopicSource(queue.getName(), groupId); + groupedTopicSource.setContextClass(source.getContextClass()); + ForkliftConsumerI consumerI = getGroupedTopic(groupedTopicSource); + return consumerI; + }) + .apply(TopicSource.class, topic -> { + GroupedTopicSource groupedTopicSource = new GroupedTopicSource(topic.getName(), groupId); + groupedTopicSource.setContextClass(source.getContextClass()); + ForkliftConsumerI consumerI = getGroupedTopic(groupedTopicSource); + return consumerI; + }) + .apply(GroupedTopicSource.class, topic -> getGroupedTopic(topic)) + .apply(RoleInputSource.class, roleSource -> { + final ForkliftConsumerI rawConsumer = getConsumerForSource(roleSource.getActionSource(this)); + return new RoleInputConsumerWrapper(rawConsumer); + }) + .elseUnsupportedError(); } public synchronized ForkliftConsumerI getGroupedTopic(GroupedTopicSource source) throws ConnectorException { if (!source.groupSpecified()) { source.overrideGroup(groupId); } - if (!source.getGroup().equals(groupId)) { //TODO actually support GroupedTopics throw new ConnectorException("Unexpected group '" + source.getGroup() + "'; only the connector group '" + groupId + "' is allowed"); } @@ -161,23 +208,22 @@ public synchronized ForkliftConsumerI getGroupedTopic(GroupedTopicSource source) if (controller != null && controller.isRunning()) { log.warn("Consumer for topic already exists under this controller's groupname. Messages will be divided amongst consumers."); } else { - controller = createController(source.getName()); + controller = createController(source); this.controllers.put(source.getName(), controller); controller.start(); } return new KafkaTopicConsumer(source.getName(), controller); } - @Override - public ForkliftConsumerI getQueue(String name) throws ConnectorException { - return getGroupedTopic(new GroupedTopicSource(name, groupId)); - } - - @Override - public ForkliftConsumerI getTopic(String name) throws ConnectorException { - return getGroupedTopic(new GroupedTopicSource(name, groupId)); - } - +// @Override +// public ForkliftConsumerI getQueue(String name) throws ConnectorException { +// return getGroupedTopic(new GroupedTopicSource(name, groupId)); +// } +// +// @Override +// public ForkliftConsumerI getTopic(String name) throws ConnectorException { +// return getGroupedTopic(new GroupedTopicSource(name, groupId)); +// } @Override public ForkliftProducerI getQueueProducer(String name) { @@ -195,8 +241,8 @@ public synchronized ForkliftProducerI getTopicProducer(String name) { @Override public ActionSource mapSource(LogicalSource source) { return source - .apply(RoleInputSource.class, roleSource -> mapRoleInputSource(roleSource)) - .get(); + .apply(RoleInputSource.class, roleSource -> mapRoleInputSource(roleSource)) + .get(); } protected GroupedTopicSource mapRoleInputSource(RoleInputSource roleSource) { diff --git a/connectors/kafka/src/main/java/forklift/message/ForkliftAvroMessageUtils.java b/connectors/kafka/src/main/java/forklift/message/ForkliftAvroMessageUtils.java new file mode 100644 index 0000000..5c64193 --- /dev/null +++ b/connectors/kafka/src/main/java/forklift/message/ForkliftAvroMessageUtils.java @@ -0,0 +1,30 @@ +package forklift.message; + +import com.fasterxml.jackson.databind.DeserializationFeature; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ArrayNode; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; +import forklift.connectors.KafkaSerializer; +import org.apache.avro.Schema; + +import java.io.IOException; + +public class ForkliftAvroMessageUtils { + + private static final ObjectMapper mapper = new ObjectMapper().registerModule(new JavaTimeModule()) + .configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false); + + public static Schema addForkliftPropertiesToSchema(Schema schema) throws IOException { + String originalJson = schema.toString(false); + JsonNode propertiesField = mapper.readTree(KafkaSerializer.SCHEMA_FIELD_VALUE_PROPERTIES); + ObjectNode schemaNode = (ObjectNode) mapper.readTree(originalJson); + ArrayNode fieldsNode = (ArrayNode) schemaNode.get("fields"); + fieldsNode.add(propertiesField); + schemaNode.set("fields", fieldsNode); + Schema.Parser parser = new Schema.Parser(); + return parser.parse(mapper.writeValueAsString(schemaNode)); + } + +} diff --git a/connectors/kafka/src/main/java/forklift/producers/KafkaForkliftProducer.java b/connectors/kafka/src/main/java/forklift/producers/KafkaForkliftProducer.java index 3cdf4c4..ef62567 100644 --- a/connectors/kafka/src/main/java/forklift/producers/KafkaForkliftProducer.java +++ b/connectors/kafka/src/main/java/forklift/producers/KafkaForkliftProducer.java @@ -1,9 +1,5 @@ package forklift.producers; -import forklift.connectors.ForkliftMessage; -import forklift.connectors.KafkaSerializer; -import forklift.message.Header; - import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.DeserializationFeature; import com.fasterxml.jackson.databind.JsonNode; @@ -11,7 +7,10 @@ import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; - +import forklift.connectors.ForkliftMessage; +import forklift.connectors.KafkaSerializer; +import forklift.message.ForkliftAvroMessageUtils; +import forklift.message.Header; import org.apache.avro.Schema; import org.apache.avro.generic.GenericData; import org.apache.avro.generic.GenericDatumReader; @@ -33,7 +32,6 @@ import java.io.ByteArrayInputStream; import java.io.ByteArrayOutputStream; import java.io.DataInputStream; -import java.io.File; import java.io.IOException; import java.io.InputStream; import java.nio.charset.Charset; @@ -109,7 +107,7 @@ public String send(ForkliftMessage message) throws ProducerException { @Override public String send(Object message) throws ProducerException { if (message instanceof SpecificRecord) { - return sendAvroMessage((SpecificRecord)message); + return sendAvroMessage((SpecificRecord) message); } else { String json; try { @@ -128,13 +126,13 @@ public String send(Map message) throws ProducerException { @Override public String send(Map headers, Map properties, ForkliftMessage message) - throws ProducerException { + throws ProducerException { throw new UnsupportedOperationException("Kafka Producer does not support headers"); } @Override public String send(Map properties, ForkliftMessage message) - throws ProducerException { + throws ProducerException { return this.sendForkliftWrappedMessage(message.getMsg(), properties); } @@ -175,7 +173,7 @@ private String sendForkliftWrappedMessage(String message, Map me avroRecord.put(KafkaSerializer.SCHEMA_FIELD_NAME_PROPERTIES, this.formatMap(appliedProperties)); ProducerRecord record = new ProducerRecord<>(topic, null, avroRecord); try { - RecordMetadata result = (RecordMetadata)kafkaProducer.send(record).get(); + RecordMetadata result = (RecordMetadata) kafkaProducer.send(record).get(); return result.topic() + "-" + result.partition() + "-" + result.offset(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); @@ -185,17 +183,6 @@ private String sendForkliftWrappedMessage(String message, Map me } } - private Schema addForkliftPropertiesToSchema(Schema schema) throws IOException { - String originalJson = schema.toString(false); - JsonNode propertiesField = mapper.readTree(KafkaSerializer.SCHEMA_FIELD_VALUE_PROPERTIES); - ObjectNode schemaNode = (ObjectNode)mapper.readTree(originalJson); - ArrayNode fieldsNode = (ArrayNode)schemaNode.get("fields"); - fieldsNode.add(propertiesField); - schemaNode.set("fields", fieldsNode); - Schema.Parser parser = new Schema.Parser(); - return parser.parse(mapper.writeValueAsString(schemaNode)); - } - private GenericRecord addForkliftPropertiesToAvroObject(SpecificRecord message) throws IOException { //Write message to json ByteArrayOutputStream outputStream = new ByteArrayOutputStream(); @@ -208,12 +195,12 @@ private GenericRecord addForkliftPropertiesToAvroObject(SpecificRecord message) //modify schema to include forklift properties Schema modifiedSchema = avroSchemaCache.get(message.getClass()); if (modifiedSchema == null) { - modifiedSchema = addForkliftPropertiesToSchema(message.getSchema()); + modifiedSchema = ForkliftAvroMessageUtils.addForkliftPropertiesToSchema(message.getSchema()); avroSchemaCache.put(message.getClass(), modifiedSchema); } //add forklift properties to json - ObjectNode messageNode = (ObjectNode)mapper.readTree(json); + ObjectNode messageNode = (ObjectNode) mapper.readTree(json); messageNode.put(KafkaSerializer.SCHEMA_FIELD_NAME_PROPERTIES, this.formatMap(this.properties)); //read modified json to avro object with modified schema @@ -227,15 +214,14 @@ private GenericRecord addForkliftPropertiesToAvroObject(SpecificRecord message) private String sendAvroMessage(SpecificRecord message) throws ProducerException { try { ProducerRecord record = null; - if(this.properties.size() > 0){ + if (this.properties.size() > 0) { GenericRecord avroRecord = addForkliftPropertiesToAvroObject(message); record = new ProducerRecord(topic, null, avroRecord); - } - else{ + } else { record = new ProducerRecord(topic, null, message); } try { - RecordMetadata result = (RecordMetadata)kafkaProducer.send(record).get(); + RecordMetadata result = (RecordMetadata) kafkaProducer.send(record).get(); return result.topic() + "-" + result.partition() + "-" + result.offset(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); diff --git a/connectors/kafka/src/main/java/forklift/serializers/ForkliftKafkaAvroDeserializer.java b/connectors/kafka/src/main/java/forklift/serializers/ForkliftKafkaAvroDeserializer.java new file mode 100644 index 0000000..8dfbeed --- /dev/null +++ b/connectors/kafka/src/main/java/forklift/serializers/ForkliftKafkaAvroDeserializer.java @@ -0,0 +1,38 @@ +package forklift.serializers; + +import io.confluent.kafka.serializers.KafkaAvroDeserializer; +import org.apache.avro.Schema; +import org.apache.kafka.common.serialization.Deserializer; + +import java.util.Map; + +/** + * Extends the {@link io.confluent.kafka.serializers.KafkaAvroDeserializer} in order to look for a reader schema + * to use while deserializing messages. This is done to override the default functionality of using the writer schema + * to deserialize GenericRecords. Using the writer schema when a reader schema is available has the potential of missing + * out on default values specified in the reader schema. A reader schema instance should be specified using the configuration + * key of "forklift.avro.reader.schema" + */ +public class ForkliftKafkaAvroDeserializer extends KafkaAvroDeserializer implements Deserializer { + + private Schema readerSchema; + + @Override + public void configure(Map configs, boolean isKey) { + readerSchema = (Schema) configs.getOrDefault("forklift.avro.reader.schema", null); + super.configure(configs, isKey); + } + + @Override + public Object deserialize(String topic, byte[] data) { + if (readerSchema != null) { + return deserialize(data, readerSchema); + } + return super.deserialize(topic, data); + } + + @Override + public void close() { + super.close(); + } +} diff --git a/connectors/kafka/src/test/java/forklift/integration/AvroMessageTests.java b/connectors/kafka/src/test/java/forklift/integration/AvroMessageTests.java index afc20a1..1cd3d2c 100644 --- a/connectors/kafka/src/test/java/forklift/integration/AvroMessageTests.java +++ b/connectors/kafka/src/test/java/forklift/integration/AvroMessageTests.java @@ -1,5 +1,7 @@ package forklift.integration; +import static org.junit.Assert.assertEquals; + import forklift.Forklift; import forklift.connectors.ConnectorException; import forklift.connectors.ForkliftMessage; @@ -9,10 +11,16 @@ import forklift.integration.server.TestServiceManager; import forklift.producers.ForkliftProducerI; import forklift.producers.ProducerException; +import forklift.schemas.AvroMessage; import forklift.schemas.StateCode; import forklift.schemas.UserRegistered; import forklift.source.decorators.Queue; - +import org.apache.avro.Schema; +import org.apache.avro.generic.GenericData; +import org.apache.avro.generic.GenericRecord; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerRecord; import org.junit.After; import org.junit.Before; import org.junit.Test; @@ -21,6 +29,12 @@ import java.util.HashMap; import java.util.Map; +import java.util.Properties; +import java.util.concurrent.Executor; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.SynchronousQueue; +import java.util.concurrent.TimeUnit; public class AvroMessageTests extends BaseIntegrationTest { @@ -42,8 +56,8 @@ public void testComplexAvroMessageWithProperty() throws ProducerException, Conne Forklift forklift = serviceManager.newManagedForkliftInstance(""); int msgCount = 10; ForkliftProducerI - producer = - forklift.getConnector().getQueueProducer("forklift-avro-topic"); + producer = + forklift.getConnector().getQueueProducer("forklift-avro-topic"); Map producerProps = new HashMap<>(); producerProps.put("Eye", "producerProperty"); producer.setProperties(producerProps); @@ -74,8 +88,8 @@ public void testComplexAvroMessageWithoutProperty() throws ProducerException, Co Forklift forklift = serviceManager.newManagedForkliftInstance(""); int msgCount = 10; ForkliftProducerI - producer = - forklift.getConnector().getQueueProducer("forklift-avro-topic"); + producer = + forklift.getConnector().getQueueProducer("forklift-avro-topic"); for (int i = 0; i < msgCount; i++) { UserRegistered registered = new UserRegistered(); registered.setFirstName("John"); @@ -88,7 +102,6 @@ public void testComplexAvroMessageWithoutProperty() throws ProducerException, Co // Shutdown the consumer after all the messages have been processed. c.setOutOfMessages((listener) -> { timeouts++; - if (sentMessageIds.equals(consumedMessageIds) || timeouts > maxTimeouts) { listener.shutdown(); } @@ -98,7 +111,61 @@ public void testComplexAvroMessageWithoutProperty() throws ProducerException, Co messageAsserts(); } - @Queue("forklift-avro-topic") + @Test + public void defaultValueInReaderSchemaAppropriatelySetWhenNotInWriterSchema() throws InterruptedException, StartupException { + Forklift forklift = serviceManager.newManagedForkliftInstance(""); + String writerSchemaWithoutColor = "{\n" + + " \"namespace\": \"forklift.schemas\",\n" + + " \"type\": \"record\",\n" + + " \"version\": 1,\n" + + " \"name\": \"AvroMessage\",\n" + + " \"fields\": [\n" + + " {\"name\": \"name\", \"type\": \"string\"}\n" + + " ]\n" + + "}"; + + Properties kafkaProps = new Properties(); + kafkaProps.put("key.serializer", io.confluent.kafka.serializers.KafkaAvroSerializer.class); + kafkaProps.put("bootstrap.servers", serviceManager.kafkaHost() + ":" + serviceManager.kafkaPort()); + kafkaProps.put("schema.registry.url", "http://" + serviceManager.schemaRegistryHost() + ":" + serviceManager.schemaRegistryPort()); + kafkaProps.put("value.serializer", io.confluent.kafka.serializers.KafkaAvroSerializer.class); + kafkaProps.put(ConsumerConfig.GROUP_ID_CONFIG, "testGroup"); + KafkaProducer producer = new KafkaProducer(kafkaProps); + Schema.Parser parser = new Schema.Parser(); + Schema schema = parser.parse(writerSchemaWithoutColor); + GenericRecord record = new GenericData.Record(schema); + record.put("name", "John"); + ProducerRecord data = new ProducerRecord("forklift-avro-default-test-topic", null, record); + producer.send(data); + + final Consumer c = new Consumer(AvroMessageConsumer.class, forklift); + ExecutorService executor = Executors.newSingleThreadExecutor(); + executor.submit(() -> c.listen()); + AvroMessage message = AvroMessageConsumer.resultHandoff.poll(5, TimeUnit.SECONDS); + c.shutdown(); + assertEquals("red", message.getColor()); + + } + + @Queue(value = "forklift-avro-default-test-topic") + public static class AvroMessageConsumer { + + public static SynchronousQueue resultHandoff = new SynchronousQueue(); + + @forklift.decorators.Message + private ForkliftMessage forkliftMessage; + + @forklift.decorators.Message + private AvroMessage value; + + @OnMessage + public void onMessage() throws InterruptedException { + resultHandoff.put(value); + consumedMessageIds.add(forkliftMessage.getId()); + } + } + + @Queue(value = "forklift-avro-topic") public static class RegisteredAvroConsumer { @forklift.decorators.Message diff --git a/connectors/kafka/src/test/java/forklift/integration/server/TestServiceManager.java b/connectors/kafka/src/test/java/forklift/integration/server/TestServiceManager.java index 6402a79..9d5d2ab 100644 --- a/connectors/kafka/src/test/java/forklift/integration/server/TestServiceManager.java +++ b/connectors/kafka/src/test/java/forklift/integration/server/TestServiceManager.java @@ -131,4 +131,22 @@ public Forklift newManagedForkliftInstance(String controllerName) throws Startup forklifts.add(forklift); return forklift; } + + + public String kafkaHost(){ + return "127.0.0.1"; + } + + public int kafkaPort(){ + return kafkaPort; + } + + public String schemaRegistryHost(){ + return "127.0.0.1"; + } + + public int schemaRegistryPort(){ + return schemaPort; + } + } diff --git a/connectors/kafka/src/test/resources/schemas/AvroMessage.avsc b/connectors/kafka/src/test/resources/schemas/AvroMessage.avsc index f8af412..8f07513 100644 --- a/connectors/kafka/src/test/resources/schemas/AvroMessage.avsc +++ b/connectors/kafka/src/test/resources/schemas/AvroMessage.avsc @@ -1,10 +1,11 @@ { "namespace": "forklift.schemas", "type": "record", - "version": 1, + "version": 2, "name": "AvroMessage", "fields": [ - {"name": "name", "type": "string"} + {"name": "name", "type": "string"}, + {"name": "color", "type": "string", "default" : "red"} ] } diff --git a/connectors/kafka/src/test/resources/schemas/ForkliftMessage.avsc b/connectors/kafka/src/test/resources/schemas/ForkliftMessage.avsc new file mode 100644 index 0000000..6bcf38e --- /dev/null +++ b/connectors/kafka/src/test/resources/schemas/ForkliftMessage.avsc @@ -0,0 +1 @@ +{"type":"record","name":"ForkliftMessage", "doc":"Non-Avro messages sent through forklift use this schema.","fields":[{"name":"forkliftValue","type":"string","default":"", "doc":"The forklift message. 3 formats are supported. 1: string value, 2: Json object,3: Map represented by key,value entries delimited with newline"},{"name":"forkliftProperties","type":"string","default":"","doc":"Properties added to support forklift interfaces. Format is key,value entries delimited with new lines"}]} diff --git a/core/build.sbt b/core/build.sbt index 376dd10..b1dae90 100644 --- a/core/build.sbt +++ b/core/build.sbt @@ -2,7 +2,7 @@ organization := "com.github.dcshock" name := "forklift" -version := "2.1" +version := "2.2" javacOptions ++= Seq("-source", "1.8") diff --git a/core/src/main/java/forklift/connectors/ForkliftConnectorI.java b/core/src/main/java/forklift/connectors/ForkliftConnectorI.java index fbc4400..e3bc54d 100644 --- a/core/src/main/java/forklift/connectors/ForkliftConnectorI.java +++ b/core/src/main/java/forklift/connectors/ForkliftConnectorI.java @@ -46,7 +46,7 @@ public interface ForkliftConnectorI extends LogicalSourceContext { * @throws ConnectorException if an error occurred interacting with the connector * @throws RuntimeException if reading from a queue is not supported */ - ForkliftConsumerI getQueue(String name) throws ConnectorException; +// ForkliftConsumerI getQueue(String name) throws ConnectorException; /** * Retrieves a {@link ForkliftConsumerI consumer} instance that reads from @@ -57,7 +57,7 @@ public interface ForkliftConnectorI extends LogicalSourceContext { * @throws ConnectorException if an error occurred interacting with the connector * @throws RuntimeException if reading from a topic is not supported */ - ForkliftConsumerI getTopic(String name) throws ConnectorException; +// ForkliftConsumerI getTopic(String name) throws ConnectorException; /** * Gives a {@link ForkliftProducerI producer} instance that writes diff --git a/core/src/main/java/forklift/source/SchemaResolver.java b/core/src/main/java/forklift/source/SchemaResolver.java new file mode 100644 index 0000000..bac9f96 --- /dev/null +++ b/core/src/main/java/forklift/source/SchemaResolver.java @@ -0,0 +1,5 @@ +package forklift.source; + +public interface SchemaResolver { + String getSchema(SourceI source); +} diff --git a/core/src/main/java/forklift/source/decorators/GroupedTopic.java b/core/src/main/java/forklift/source/decorators/GroupedTopic.java index f96124e..af1e625 100644 --- a/core/src/main/java/forklift/source/decorators/GroupedTopic.java +++ b/core/src/main/java/forklift/source/decorators/GroupedTopic.java @@ -1,5 +1,7 @@ package forklift.source.decorators; +import forklift.source.SchemaResolver; +import forklift.source.SourceI; import forklift.source.SourceType; import forklift.source.sources.GroupedTopicSource; @@ -30,4 +32,6 @@ * If empty, some consumer group name should be generated. */ String group() default ""; + +// Class schemaResolver() default Topic.None.class; } diff --git a/core/src/main/java/forklift/source/decorators/Queue.java b/core/src/main/java/forklift/source/decorators/Queue.java index 9358cf1..9566181 100644 --- a/core/src/main/java/forklift/source/decorators/Queue.java +++ b/core/src/main/java/forklift/source/decorators/Queue.java @@ -1,11 +1,11 @@ package forklift.source.decorators; +import forklift.source.SchemaResolver; import forklift.source.SourceType; import forklift.source.sources.QueueSource; import java.lang.annotation.Documented; import java.lang.annotation.ElementType; -import java.lang.annotation.Inherited; import java.lang.annotation.Repeatable; import java.lang.annotation.Retention; import java.lang.annotation.RetentionPolicy; @@ -21,4 +21,7 @@ @Target({ElementType.TYPE}) public @interface Queue { String value(); + + //Class schemaResolver() default Topic.None.class; + } diff --git a/core/src/main/java/forklift/source/decorators/Topic.java b/core/src/main/java/forklift/source/decorators/Topic.java index edff854..dbd03bb 100644 --- a/core/src/main/java/forklift/source/decorators/Topic.java +++ b/core/src/main/java/forklift/source/decorators/Topic.java @@ -1,5 +1,7 @@ package forklift.source.decorators; +import forklift.source.SchemaResolver; +import forklift.source.SourceI; import forklift.source.SourceType; import forklift.source.sources.TopicSource; @@ -21,4 +23,19 @@ @Target({ElementType.TYPE}) public @interface Topic { String value(); +// Class schemaResolver() default None.class; +// +// //because you can't use null as a default value in an annotation. +// class None implements SchemaResolver { +// private static final long serialVersionUID = 1L; +// +// private None() { +// } +// +// @Override +// public String getSchema(SourceI source) { +// return null; +// } +// } + } From 7c76408610819591316b20410aa02e7e50bfd793 Mon Sep 17 00:00:00 2001 From: Andrew Frieze Date: Mon, 20 Nov 2017 15:57:43 -0700 Subject: [PATCH 2/2] Code cleanup and general readability enhancements --- .../forklift/connectors/KafkaConnector.java | 25 +++------------ .../message/ForkliftAvroMessageUtils.java | 7 ++++ .../ForkliftKafkaAvroDeserializer.java | 5 +-- .../connectors/ForkliftConnectorI.java | 32 +++---------------- .../source/decorators/GroupedTopic.java | 5 +-- .../forklift/source/decorators/Queue.java | 2 -- .../forklift/source/decorators/Topic.java | 18 ----------- 7 files changed, 20 insertions(+), 74 deletions(-) diff --git a/connectors/kafka/src/main/java/forklift/connectors/KafkaConnector.java b/connectors/kafka/src/main/java/forklift/connectors/KafkaConnector.java index 112a437..8309a4b 100644 --- a/connectors/kafka/src/main/java/forklift/connectors/KafkaConnector.java +++ b/connectors/kafka/src/main/java/forklift/connectors/KafkaConnector.java @@ -1,11 +1,5 @@ package forklift.connectors; -import com.fasterxml.jackson.databind.DeserializationFeature; -import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.ObjectMapper; -import com.fasterxml.jackson.databind.node.ArrayNode; -import com.fasterxml.jackson.databind.node.ObjectNode; -import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; import forklift.consumer.ForkliftConsumerI; import forklift.consumer.KafkaTopicConsumer; import forklift.consumer.wrapper.RoleInputConsumerWrapper; @@ -121,7 +115,7 @@ private KafkaController createController(GroupedTopicSource source) { ForkliftKafkaAvroDeserializer deserializer = new ForkliftKafkaAvroDeserializer(); Schema readerSchema = null; try { - readerSchema = produceReaderSchema(source); + readerSchema = readerSchemaFromSource(source); } catch (Exception e) { log.error("Unable to generate reader schema, falling back to writer schema. Defaults may be lost"); } @@ -133,7 +127,7 @@ private KafkaController createController(GroupedTopicSource source) { props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, io.confluent.kafka.serializers.KafkaAvroDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ForkliftKafkaAvroDeserializer.class); if (readerSchema != null) { - props.put("forklift.avro.reader.schema", readerSchema); + props.put(ForkliftKafkaAvroDeserializer.READER_SCHEMA_CONFIG, readerSchema); } props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 200); @@ -144,8 +138,7 @@ private KafkaController createController(GroupedTopicSource source) { return new KafkaController(kafkaConsumer, new MessageStream(), source.getName()); } - private Schema produceReaderSchema(SourceI source) throws NoSuchMethodException, InvocationTargetException, IllegalAccessException, IOException { - Set fields = new HashSet<>(); + private Schema readerSchemaFromSource(SourceI source) throws NoSuchMethodException, InvocationTargetException, IllegalAccessException, IOException { for (Field field : source.getContextClass().getDeclaredFields()) { if (field.isAnnotationPresent(Message.class) && SpecificRecord.class.isAssignableFrom(field.getType())) { Schema schema = (Schema) field.getType().getMethod("getClassSchema").invoke(null); @@ -214,17 +207,7 @@ public synchronized ForkliftConsumerI getGroupedTopic(GroupedTopicSource source) } return new KafkaTopicConsumer(source.getName(), controller); } - -// @Override -// public ForkliftConsumerI getQueue(String name) throws ConnectorException { -// return getGroupedTopic(new GroupedTopicSource(name, groupId)); -// } -// -// @Override -// public ForkliftConsumerI getTopic(String name) throws ConnectorException { -// return getGroupedTopic(new GroupedTopicSource(name, groupId)); -// } - + @Override public ForkliftProducerI getQueueProducer(String name) { return getTopicProducer(name); diff --git a/connectors/kafka/src/main/java/forklift/message/ForkliftAvroMessageUtils.java b/connectors/kafka/src/main/java/forklift/message/ForkliftAvroMessageUtils.java index 5c64193..03cfeb3 100644 --- a/connectors/kafka/src/main/java/forklift/message/ForkliftAvroMessageUtils.java +++ b/connectors/kafka/src/main/java/forklift/message/ForkliftAvroMessageUtils.java @@ -16,6 +16,13 @@ public class ForkliftAvroMessageUtils { private static final ObjectMapper mapper = new ObjectMapper().registerModule(new JavaTimeModule()) .configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false); + /** + * Adds the forkliftProperties field to the passed in Schema. + * + * @param schema the schema which forklift properties should be added + * @return Schema with forlift properties added + * @throws IOException + */ public static Schema addForkliftPropertiesToSchema(Schema schema) throws IOException { String originalJson = schema.toString(false); JsonNode propertiesField = mapper.readTree(KafkaSerializer.SCHEMA_FIELD_VALUE_PROPERTIES); diff --git a/connectors/kafka/src/main/java/forklift/serializers/ForkliftKafkaAvroDeserializer.java b/connectors/kafka/src/main/java/forklift/serializers/ForkliftKafkaAvroDeserializer.java index 8dfbeed..6a8ba45 100644 --- a/connectors/kafka/src/main/java/forklift/serializers/ForkliftKafkaAvroDeserializer.java +++ b/connectors/kafka/src/main/java/forklift/serializers/ForkliftKafkaAvroDeserializer.java @@ -11,15 +11,16 @@ * to use while deserializing messages. This is done to override the default functionality of using the writer schema * to deserialize GenericRecords. Using the writer schema when a reader schema is available has the potential of missing * out on default values specified in the reader schema. A reader schema instance should be specified using the configuration - * key of "forklift.avro.reader.schema" + * key of {@link #READER_SCHEMA_CONFIG} */ public class ForkliftKafkaAvroDeserializer extends KafkaAvroDeserializer implements Deserializer { + public static final String READER_SCHEMA_CONFIG = "forklift.avro.reader.schema"; private Schema readerSchema; @Override public void configure(Map configs, boolean isKey) { - readerSchema = (Schema) configs.getOrDefault("forklift.avro.reader.schema", null); + readerSchema = (Schema) configs.getOrDefault(READER_SCHEMA_CONFIG, null); super.configure(configs, isKey); } diff --git a/core/src/main/java/forklift/connectors/ForkliftConnectorI.java b/core/src/main/java/forklift/connectors/ForkliftConnectorI.java index e3bc54d..3bccbd3 100644 --- a/core/src/main/java/forklift/connectors/ForkliftConnectorI.java +++ b/core/src/main/java/forklift/connectors/ForkliftConnectorI.java @@ -2,8 +2,8 @@ import forklift.consumer.ForkliftConsumerI; import forklift.producers.ForkliftProducerI; -import forklift.source.SourceI; import forklift.source.LogicalSourceContext; +import forklift.source.SourceI; /** * An entity that manages consuming and producing to a particular message bus. @@ -14,7 +14,7 @@ public interface ForkliftConnectorI extends LogicalSourceContext { * connector methods to be expected to behave properly. * * @throws ConnectorException if there was a problem initializing the state - * of this connector or making a connection + * of this connector or making a connection */ void start() throws ConnectorException; @@ -22,7 +22,7 @@ public interface ForkliftConnectorI extends LogicalSourceContext { * Stops the given connector, invalidating it's current state. * * @throws ConnectorException if there was a problem un-initializing the state - * of this connector + * of this connector */ void stop() throws ConnectorException; @@ -33,32 +33,10 @@ public interface ForkliftConnectorI extends LogicalSourceContext { * @param source the source to read from * @return a consumer that reads from the given source * @throws ConnectorException if an error occurred interacting with the connector - * @throws RuntimeException if reading from the given source is not supported + * @throws RuntimeException if reading from the given source is not supported */ ForkliftConsumerI getConsumerForSource(SourceI source) throws ConnectorException; - /** - * Retrieves the {@link ForkliftConsumerI consumer} instance that reads from - * the queue with the given name. - * - * @param name the name of the topic to read from - * @return a consumer that reads from the given topic - * @throws ConnectorException if an error occurred interacting with the connector - * @throws RuntimeException if reading from a queue is not supported - */ -// ForkliftConsumerI getQueue(String name) throws ConnectorException; - - /** - * Retrieves a {@link ForkliftConsumerI consumer} instance that reads from - * the topic with the given name. - * - * @param name the name of the topic to read from - * @return a consumer that reads from the given topic - * @throws ConnectorException if an error occurred interacting with the connector - * @throws RuntimeException if reading from a topic is not supported - */ -// ForkliftConsumerI getTopic(String name) throws ConnectorException; - /** * Gives a {@link ForkliftProducerI producer} instance that writes * to the given queue. @@ -83,7 +61,7 @@ public interface ForkliftConnectorI extends LogicalSourceContext { * need to be sent by an external application. * * @return the serializer to use for serializing messages to this connector, - * or null if there is no explicit serialization method for this connector + * or null if there is no explicit serialization method for this connector */ default ForkliftSerializer getDefaultSerializer() { return null; diff --git a/core/src/main/java/forklift/source/decorators/GroupedTopic.java b/core/src/main/java/forklift/source/decorators/GroupedTopic.java index af1e625..f83d214 100644 --- a/core/src/main/java/forklift/source/decorators/GroupedTopic.java +++ b/core/src/main/java/forklift/source/decorators/GroupedTopic.java @@ -1,7 +1,5 @@ package forklift.source.decorators; -import forklift.source.SchemaResolver; -import forklift.source.SourceI; import forklift.source.SourceType; import forklift.source.sources.GroupedTopicSource; @@ -28,10 +26,9 @@ /** * The name of the consumer group reading this topic. - * + *

* If empty, some consumer group name should be generated. */ String group() default ""; -// Class schemaResolver() default Topic.None.class; } diff --git a/core/src/main/java/forklift/source/decorators/Queue.java b/core/src/main/java/forklift/source/decorators/Queue.java index 9566181..56b146d 100644 --- a/core/src/main/java/forklift/source/decorators/Queue.java +++ b/core/src/main/java/forklift/source/decorators/Queue.java @@ -22,6 +22,4 @@ public @interface Queue { String value(); - //Class schemaResolver() default Topic.None.class; - } diff --git a/core/src/main/java/forklift/source/decorators/Topic.java b/core/src/main/java/forklift/source/decorators/Topic.java index dbd03bb..daa37b2 100644 --- a/core/src/main/java/forklift/source/decorators/Topic.java +++ b/core/src/main/java/forklift/source/decorators/Topic.java @@ -1,13 +1,10 @@ package forklift.source.decorators; -import forklift.source.SchemaResolver; -import forklift.source.SourceI; import forklift.source.SourceType; import forklift.source.sources.TopicSource; import java.lang.annotation.Documented; import java.lang.annotation.ElementType; -import java.lang.annotation.Inherited; import java.lang.annotation.Repeatable; import java.lang.annotation.Retention; import java.lang.annotation.RetentionPolicy; @@ -23,19 +20,4 @@ @Target({ElementType.TYPE}) public @interface Topic { String value(); -// Class schemaResolver() default None.class; -// -// //because you can't use null as a default value in an annotation. -// class None implements SchemaResolver { -// private static final long serialVersionUID = 1L; -// -// private None() { -// } -// -// @Override -// public String getSchema(SourceI source) { -// return null; -// } -// } - }