This is an automated email from the ASF dual-hosted git repository.

ppkarwasz pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/logging-flume-legacy.git

commit 99b23fd3e1cc49ee84ec7ef1723688e96aed8649
Author: Mike Percy <[email protected]>
AuthorDate: Tue Jun 18 11:17:39 2013 -0700

    FLUME-2010. Support Avro records in Log4jAppender and the HDFS Sink.
    
    (Tom White via Mike Percy)
---
 .../flume/clients/log4jappender/Log4jAppender.java |  73 +++++++-
 .../clients/log4jappender/Log4jAvroHeaders.java    |   4 +-
 .../log4jappender/TestLog4jAppenderWithAvro.java   | 195 +++++++++++++++++++++
 .../flume-log4jtest-avro-generic.properties        |  21 +++
 .../flume-log4jtest-avro-reflect.properties        |  21 +++
 .../src/test/resources/myrecord.avsc               |   1 +
 6 files changed, 306 insertions(+), 9 deletions(-)

diff --git 
a/flume-ng-log4jappender/src/main/java/org/apache/flume/clients/log4jappender/Log4jAppender.java
 
b/flume-ng-log4jappender/src/main/java/org/apache/flume/clients/log4jappender/Log4jAppender.java
index 532b761..b07b189 100644
--- 
a/flume-ng-log4jappender/src/main/java/org/apache/flume/clients/log4jappender/Log4jAppender.java
+++ 
b/flume-ng-log4jappender/src/main/java/org/apache/flume/clients/log4jappender/Log4jAppender.java
@@ -18,11 +18,21 @@
  */
 package org.apache.flume.clients.log4jappender;
 
+import java.io.ByteArrayOutputStream;
+import java.io.IOException;
 import java.nio.charset.Charset;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Properties;
 
+import org.apache.avro.Schema;
+import org.apache.avro.generic.GenericRecord;
+import org.apache.avro.io.BinaryEncoder;
+import org.apache.avro.io.DatumWriter;
+import org.apache.avro.io.EncoderFactory;
+import org.apache.avro.reflect.ReflectData;
+import org.apache.avro.reflect.ReflectDatumWriter;
+import org.apache.avro.specific.SpecificRecord;
 import org.apache.flume.Event;
 import org.apache.flume.EventDeliveryException;
 import org.apache.flume.FlumeException;
@@ -67,6 +77,8 @@ public class Log4jAppender extends AppenderSkeleton {
   private boolean unsafeMode = false;
   private long timeout = RpcClientConfigurationConstants
     .DEFAULT_REQUEST_TIMEOUT_MILLIS;
+  private boolean avroReflectionEnabled;
+  private String avroSchemaUrl;
 
   RpcClient rpcClient = null;
 
@@ -130,18 +142,23 @@ public class Log4jAppender extends AppenderSkeleton {
     //Log4jAvroHeaders.LOG_LEVEL.toString()))
     hdrs.put(Log4jAvroHeaders.LOG_LEVEL.toString(),
         String.valueOf(event.getLevel().toInt()));
-    hdrs.put(Log4jAvroHeaders.MESSAGE_ENCODING.toString(), "UTF8");
 
-    String message = null;
-    if(this.layout != null) {
-      message = this.layout.format(event);
+    Event flumeEvent;
+    Object message = event.getMessage();
+    if (message instanceof GenericRecord) {
+      GenericRecord record = (GenericRecord) message;
+      populateAvroHeaders(hdrs, record.getSchema(), message);
+      flumeEvent = EventBuilder.withBody(serialize(record, 
record.getSchema()), hdrs);
+    } else if (message instanceof SpecificRecord || avroReflectionEnabled) {
+      Schema schema = ReflectData.get().getSchema(message.getClass());
+      populateAvroHeaders(hdrs, schema, message);
+      flumeEvent = EventBuilder.withBody(serialize(message, schema), hdrs);
     } else {
-      message = event.getMessage().toString();
+      hdrs.put(Log4jAvroHeaders.MESSAGE_ENCODING.toString(), "UTF8");
+      String msg = layout != null ? layout.format(event) : message.toString();
+      flumeEvent = EventBuilder.withBody(msg, Charset.forName("UTF8"), hdrs);
     }
 
-    Event flumeEvent = EventBuilder.withBody(
-        message, Charset.forName("UTF8"), hdrs);
-
     try {
       rpcClient.append(flumeEvent);
     } catch (EventDeliveryException e) {
@@ -154,6 +171,39 @@ public class Log4jAppender extends AppenderSkeleton {
     }
   }
 
+  private Schema schema;
+  private ByteArrayOutputStream out;
+  private DatumWriter<Object> writer;
+  private BinaryEncoder encoder;
+
+  protected void populateAvroHeaders(Map<String, String> hdrs, Schema schema,
+      Object message) {
+    if (avroSchemaUrl != null) {
+      hdrs.put(Log4jAvroHeaders.AVRO_SCHEMA_URL.toString(), avroSchemaUrl);
+      return;
+    }
+    LogLog.warn("Cannot find ID for schema. Adding header for schema, " +
+        "which may be inefficient. Consider setting up an Avro Schema Cache.");
+    hdrs.put(Log4jAvroHeaders.AVRO_SCHEMA_LITERAL.toString(), 
schema.toString());
+  }
+
+  private byte[] serialize(Object datum, Schema datumSchema) throws 
FlumeException {
+    if (schema == null || !datumSchema.equals(schema)) {
+      schema = datumSchema;
+      out = new ByteArrayOutputStream();
+      writer = new ReflectDatumWriter<Object>(schema);
+      encoder = EncoderFactory.get().binaryEncoder(out, null);
+    }
+    out.reset();
+    try {
+      writer.write(datum, encoder);
+      encoder.flush();
+      return out.toByteArray();
+    } catch (IOException e) {
+      throw new FlumeException(e);
+    }
+  }
+
   //This function should be synchronized to make sure one thread
   //does not close an appender another thread is using, and hence risking
   //a null pointer exception.
@@ -229,6 +279,13 @@ public class Log4jAppender extends AppenderSkeleton {
     return this.timeout;
   }
 
+  public void setAvroReflectionEnabled(boolean avroReflectionEnabled) {
+    this.avroReflectionEnabled = avroReflectionEnabled;
+  }
+
+  public void setAvroSchemaUrl(String avroSchemaUrl) {
+    this.avroSchemaUrl = avroSchemaUrl;
+  }
 
   /**
    * Activate the options set using <tt>setPort()</tt>
diff --git 
a/flume-ng-log4jappender/src/main/java/org/apache/flume/clients/log4jappender/Log4jAvroHeaders.java
 
b/flume-ng-log4jappender/src/main/java/org/apache/flume/clients/log4jappender/Log4jAvroHeaders.java
index a6216c3..08a7203 100644
--- 
a/flume-ng-log4jappender/src/main/java/org/apache/flume/clients/log4jappender/Log4jAvroHeaders.java
+++ 
b/flume-ng-log4jappender/src/main/java/org/apache/flume/clients/log4jappender/Log4jAvroHeaders.java
@@ -23,7 +23,9 @@ public enum Log4jAvroHeaders {
   LOGGER_NAME("flume.client.log4j.logger.name"),
   LOG_LEVEL("flume.client.log4j.log.level"),
   MESSAGE_ENCODING("flume.client.log4j.message.encoding"),
-  TIMESTAMP("flume.client.log4j.timestamp");
+  TIMESTAMP("flume.client.log4j.timestamp"),
+  AVRO_SCHEMA_LITERAL("flume.avro.schema.literal"),
+  AVRO_SCHEMA_URL("flume.avro.schema.url");
 
   private String headerName;
   private Log4jAvroHeaders(String headerName){
diff --git 
a/flume-ng-log4jappender/src/test/java/org/apache/flume/clients/log4jappender/TestLog4jAppenderWithAvro.java
 
b/flume-ng-log4jappender/src/test/java/org/apache/flume/clients/log4jappender/TestLog4jAppenderWithAvro.java
new file mode 100644
index 0000000..5899c62
--- /dev/null
+++ 
b/flume-ng-log4jappender/src/test/java/org/apache/flume/clients/log4jappender/TestLog4jAppenderWithAvro.java
@@ -0,0 +1,195 @@
+/*
+ * 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.flume.clients.log4jappender;
+
+import com.google.common.io.Files;
+import com.google.common.io.Resources;
+import java.io.File;
+import java.io.FileReader;
+import java.io.IOException;
+import java.net.URL;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.Properties;
+import junit.framework.Assert;
+import org.apache.avro.Schema;
+import org.apache.avro.generic.GenericDatumReader;
+import org.apache.avro.generic.GenericRecord;
+import org.apache.avro.generic.GenericRecordBuilder;
+import org.apache.avro.io.BinaryDecoder;
+import org.apache.avro.io.DecoderFactory;
+import org.apache.avro.reflect.ReflectData;
+import org.apache.avro.reflect.ReflectDatumReader;
+import org.apache.flume.Channel;
+import org.apache.flume.ChannelSelector;
+import org.apache.flume.Context;
+import org.apache.flume.Event;
+import org.apache.flume.Transaction;
+import org.apache.flume.channel.ChannelProcessor;
+import org.apache.flume.channel.MemoryChannel;
+import org.apache.flume.channel.ReplicatingChannelSelector;
+import org.apache.flume.conf.Configurables;
+import org.apache.flume.source.AvroSource;
+import org.apache.log4j.LogManager;
+import org.apache.log4j.Logger;
+import org.apache.log4j.PropertyConfigurator;
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Test;
+
+public class TestLog4jAppenderWithAvro {
+  private AvroSource source;
+  private Channel ch;
+  private Properties props;
+
+  @Before
+  public void setUp() throws Exception {
+    URL schemaUrl = getClass().getClassLoader().getResource("myrecord.avsc");
+    Files.copy(Resources.newInputStreamSupplier(schemaUrl),
+        new File("/tmp/myrecord.avsc"));
+
+    int port = 25430;
+    source = new AvroSource();
+    ch = new MemoryChannel();
+    Configurables.configure(ch, new Context());
+
+    Context context = new Context();
+    context.put("port", String.valueOf(port));
+    context.put("bind", "localhost");
+    Configurables.configure(source, context);
+
+    List<Channel> channels = new ArrayList<Channel>();
+    channels.add(ch);
+
+    ChannelSelector rcs = new ReplicatingChannelSelector();
+    rcs.setChannels(channels);
+
+    source.setChannelProcessor(new ChannelProcessor(rcs));
+
+    source.start();
+  }
+
+  private void loadProperties(String file) throws IOException {
+    File TESTFILE = new File(
+        TestLog4jAppenderWithAvro.class.getClassLoader()
+            .getResource(file).getFile());
+    FileReader reader = new FileReader(TESTFILE);
+    props = new Properties();
+    props.load(reader);
+    reader.close();
+  }
+
+  @Test
+  public void testAvroGeneric() throws IOException {
+    loadProperties("flume-log4jtest-avro-generic.properties");
+    PropertyConfigurator.configure(props);
+    Logger logger = LogManager.getLogger(TestLog4jAppenderWithAvro.class);
+    String msg = "This is log message number " + String.valueOf(0);
+
+    Schema schema = new Schema.Parser().parse(
+        getClass().getClassLoader().getResource("myrecord.avsc").openStream());
+    GenericRecordBuilder builder = new GenericRecordBuilder(schema);
+    GenericRecord record = builder.set("message", msg).build();
+
+    logger.info(record);
+
+    Transaction transaction = ch.getTransaction();
+    transaction.begin();
+    Event event = ch.take();
+    Assert.assertNotNull(event);
+
+    GenericDatumReader<GenericRecord> reader = new 
GenericDatumReader<GenericRecord>(schema);
+    BinaryDecoder decoder = 
DecoderFactory.get().binaryDecoder(event.getBody(), null);
+    GenericRecord recordFromEvent = reader.read(null, decoder);
+    Assert.assertEquals(msg, recordFromEvent.get("message").toString());
+
+    Map<String, String> hdrs = event.getHeaders();
+
+    Assert.assertNull(hdrs.get(Log4jAvroHeaders.MESSAGE_ENCODING.toString()));
+
+    Assert.assertEquals("Schema URL should be set",
+        "file:///tmp/myrecord.avsc", 
hdrs.get(Log4jAvroHeaders.AVRO_SCHEMA_URL.toString
+        ()));
+    Assert.assertNull("Schema string should not be set",
+        hdrs.get(Log4jAvroHeaders.AVRO_SCHEMA_LITERAL.toString()));
+
+    transaction.commit();
+    transaction.close();
+
+  }
+
+  @Test
+  public void testAvroReflect() throws IOException {
+    loadProperties("flume-log4jtest-avro-reflect.properties");
+    PropertyConfigurator.configure(props);
+    Logger logger = LogManager.getLogger(TestLog4jAppenderWithAvro.class);
+    String msg = "This is log message number " + String.valueOf(0);
+
+    AppEvent appEvent = new AppEvent();
+    appEvent.setMessage(msg);
+
+    logger.info(appEvent);
+
+    Transaction transaction = ch.getTransaction();
+    transaction.begin();
+    Event event = ch.take();
+    Assert.assertNotNull(event);
+
+    Schema schema = ReflectData.get().getSchema(appEvent.getClass());
+
+    ReflectDatumReader<AppEvent> reader = new 
ReflectDatumReader<AppEvent>(AppEvent.class);
+    BinaryDecoder decoder = 
DecoderFactory.get().binaryDecoder(event.getBody(), null);
+    AppEvent recordFromEvent = reader.read(null, decoder);
+    Assert.assertEquals(msg, recordFromEvent.getMessage());
+
+    Map<String, String> hdrs = event.getHeaders();
+
+    Assert.assertNull(hdrs.get(Log4jAvroHeaders.MESSAGE_ENCODING.toString()));
+
+    Assert.assertNull("Schema URL should not be set",
+        hdrs.get(Log4jAvroHeaders.AVRO_SCHEMA_URL.toString()));
+    Assert.assertEquals("Schema string should be set", schema.toString(),
+        hdrs.get(Log4jAvroHeaders.AVRO_SCHEMA_LITERAL.toString()));
+
+    transaction.commit();
+    transaction.close();
+
+  }
+
+  @After
+  public void cleanUp(){
+    source.stop();
+    ch.stop();
+    props.clear();
+  }
+
+  public static class AppEvent {
+    private String message;
+
+    public String getMessage() {
+      return message;
+    }
+
+    public void setMessage(String message) {
+      this.message = message;
+    }
+  }
+
+}
diff --git 
a/flume-ng-log4jappender/src/test/resources/flume-log4jtest-avro-generic.properties
 
b/flume-ng-log4jappender/src/test/resources/flume-log4jtest-avro-generic.properties
new file mode 100644
index 0000000..ffdab8b
--- /dev/null
+++ 
b/flume-ng-log4jappender/src/test/resources/flume-log4jtest-avro-generic.properties
@@ -0,0 +1,21 @@
+# 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.
+log4j.appender.out2 = org.apache.flume.clients.log4jappender.Log4jAppender
+log4j.appender.out2.Port = 25430
+log4j.appender.out2.Hostname = localhost
+log4j.appender.out2.AvroSchemaUrl = file:///tmp/myrecord.avsc
+log4j.logger.org.apache.flume.clients.log4jappender = DEBUG,out2
\ No newline at end of file
diff --git 
a/flume-ng-log4jappender/src/test/resources/flume-log4jtest-avro-reflect.properties
 
b/flume-ng-log4jappender/src/test/resources/flume-log4jtest-avro-reflect.properties
new file mode 100644
index 0000000..b50ffcc
--- /dev/null
+++ 
b/flume-ng-log4jappender/src/test/resources/flume-log4jtest-avro-reflect.properties
@@ -0,0 +1,21 @@
+# 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.
+log4j.appender.out2 = org.apache.flume.clients.log4jappender.Log4jAppender
+log4j.appender.out2.Port = 25430
+log4j.appender.out2.Hostname = localhost
+log4j.appender.out2.AvroReflectionEnabled = true
+log4j.logger.org.apache.flume.clients.log4jappender = DEBUG,out2
\ No newline at end of file
diff --git a/flume-ng-log4jappender/src/test/resources/myrecord.avsc 
b/flume-ng-log4jappender/src/test/resources/myrecord.avsc
new file mode 100644
index 0000000..54130a3
--- /dev/null
+++ b/flume-ng-log4jappender/src/test/resources/myrecord.avsc
@@ -0,0 +1 @@
+{"type":"record","name":"myrecord","fields":[{"name":"message","type":"string"}]}
\ No newline at end of file

Reply via email to