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
