View Javadoc
1   /*
2    * Licensed to the Apache Software Foundation (ASF) under one or more
3    * contributor license agreements. See the NOTICE file distributed with
4    * this work for additional information regarding copyright ownership.
5    * The ASF licenses this file to You under the Apache license, Version 2.0
6    * (the "License"); you may not use this file except in compliance with
7    * the License. You may obtain a copy of the License at
8    *
9    *      http://www.apache.org/licenses/LICENSE-2.0
10   *
11   * Unless required by applicable law or agreed to in writing, software
12   * distributed under the License is distributed on an "AS IS" BASIS,
13   * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14   * See the license for the specific language governing permissions and
15   * limitations under the license.
16   */
17  package org.apache.logging.log4j.cassandra;
18  
19  import java.net.InetSocketAddress;
20  import java.util.ArrayList;
21  import java.util.Date;
22  import java.util.List;
23  
24  import com.datastax.driver.core.BatchStatement;
25  import com.datastax.driver.core.BoundStatement;
26  import com.datastax.driver.core.Cluster;
27  import com.datastax.driver.core.PreparedStatement;
28  import com.datastax.driver.core.Session;
29  import org.apache.logging.log4j.core.LogEvent;
30  import org.apache.logging.log4j.core.appender.ManagerFactory;
31  import org.apache.logging.log4j.core.appender.db.AbstractDatabaseManager;
32  import org.apache.logging.log4j.core.appender.db.ColumnMapping;
33  import org.apache.logging.log4j.core.config.plugins.convert.DateTypeConverter;
34  import org.apache.logging.log4j.core.config.plugins.convert.TypeConverters;
35  import org.apache.logging.log4j.core.net.SocketAddress;
36  import org.apache.logging.log4j.spi.ThreadContextMap;
37  import org.apache.logging.log4j.spi.ThreadContextStack;
38  import org.apache.logging.log4j.util.ReadOnlyStringMap;
39  import org.apache.logging.log4j.util.Strings;
40  
41  /**
42   * Manager for a Cassandra appender instance.
43   */
44  public class CassandraManager extends AbstractDatabaseManager {
45  
46      private static final int DEFAULT_PORT = 9042;
47  
48      private final Cluster cluster;
49      private final String keyspace;
50      private final String insertQueryTemplate;
51      private final List<ColumnMapping> columnMappings;
52      private final BatchStatement batchStatement;
53      // re-usable argument binding array
54      private final Object[] values;
55  
56      private Session session;
57      private PreparedStatement preparedStatement;
58  
59      private CassandraManager(final String name, final int bufferSize, final Cluster cluster,
60                               final String keyspace, final String insertQueryTemplate,
61                               final List<ColumnMapping> columnMappings, final BatchStatement batchStatement) {
62          super(name, bufferSize);
63          this.cluster = cluster;
64          this.keyspace = keyspace;
65          this.insertQueryTemplate = insertQueryTemplate;
66          this.columnMappings = columnMappings;
67          this.batchStatement = batchStatement;
68          this.values = new Object[columnMappings.size()];
69      }
70  
71      @Override
72      protected void startupInternal() throws Exception {
73          session = cluster.connect(keyspace);
74          preparedStatement = session.prepare(insertQueryTemplate);
75      }
76  
77      @Override
78      protected boolean shutdownInternal() throws Exception {
79          session.close();
80          cluster.close();
81          return true;
82      }
83  
84      @Override
85      protected void connectAndStart() {
86          // a Session automatically manages connections for us
87      }
88  
89      @Override
90      protected void writeInternal(final LogEvent event) {
91          for (int i = 0; i < columnMappings.size(); i++) {
92              final ColumnMapping columnMapping = columnMappings.get(i);
93              if (ThreadContextMap.class.isAssignableFrom(columnMapping.getType())
94                  || ReadOnlyStringMap.class.isAssignableFrom(columnMapping.getType())) {
95                  values[i] = event.getContextData().toMap();
96              } else if (ThreadContextStack.class.isAssignableFrom(columnMapping.getType())) {
97                  values[i] = event.getContextStack().asList();
98              } else if (Date.class.isAssignableFrom(columnMapping.getType())) {
99                  values[i] = DateTypeConverter.fromMillis(event.getTimeMillis(), columnMapping.getType().asSubclass(Date.class));
100             } else {
101                 values[i] = TypeConverters.convert(columnMapping.getLayout().toSerializable(event),
102                     columnMapping.getType(), null);
103             }
104         }
105         final BoundStatement boundStatement = preparedStatement.bind(values);
106         if (batchStatement == null) {
107             session.execute(boundStatement);
108         } else {
109             batchStatement.add(boundStatement);
110         }
111     }
112 
113     @Override
114     protected boolean commitAndClose() {
115         if (batchStatement != null) {
116             session.execute(batchStatement);
117         }
118         return true;
119     }
120 
121     public static CassandraManager getManager(final String name, final SocketAddress[] contactPoints,
122                                               final ColumnMapping[] columns, final boolean useTls,
123                                               final String clusterName, final String keyspace, final String table,
124                                               final String username, final String password,
125                                               final boolean useClockForTimestampGenerator, final int bufferSize,
126                                               final boolean batched, final BatchStatement.Type batchType) {
127         return getManager(name,
128             new FactoryData(contactPoints, columns, useTls, clusterName, keyspace, table, username, password,
129                 useClockForTimestampGenerator, bufferSize, batched, batchType), CassandraManagerFactory.INSTANCE);
130     }
131 
132     private static class CassandraManagerFactory implements ManagerFactory<CassandraManager, FactoryData> {
133 
134         private static final CassandraManagerFactory INSTANCE = new CassandraManagerFactory();
135 
136         @Override
137         public CassandraManager createManager(final String name, final FactoryData data) {
138             final Cluster.Builder builder = Cluster.builder()
139                 .addContactPointsWithPorts(data.contactPoints)
140                 .withClusterName(data.clusterName);
141             if (data.useTls) {
142                 builder.withSSL();
143             }
144             if (Strings.isNotBlank(data.username)) {
145                 builder.withCredentials(data.username, data.password);
146             }
147             if (data.useClockForTimestampGenerator) {
148                 builder.withTimestampGenerator(new ClockTimestampGenerator());
149             }
150             final Cluster cluster = builder.build();
151 
152             final StringBuilder sb = new StringBuilder("INSERT INTO ").append(data.table).append(" (");
153             for (final ColumnMapping column : data.columns) {
154                 sb.append(column.getName()).append(',');
155             }
156             sb.setCharAt(sb.length() - 1, ')');
157             sb.append(" VALUES (");
158             final List<ColumnMapping> columnMappings = new ArrayList<>(data.columns.length);
159             for (final ColumnMapping column : data.columns) {
160                 if (Strings.isNotEmpty(column.getLiteralValue())) {
161                     sb.append(column.getLiteralValue());
162                 } else {
163                     sb.append('?');
164                     columnMappings.add(column);
165                 }
166                 sb.append(',');
167             }
168             sb.setCharAt(sb.length() - 1, ')');
169             final String insertQueryTemplate = sb.toString();
170             LOGGER.debug("Using CQL for appender {}: {}", name, insertQueryTemplate);
171             return new CassandraManager(name, data.getBufferSize(), cluster, data.keyspace, insertQueryTemplate,
172                 columnMappings, data.batched ? new BatchStatement(data.batchType) : null);
173         }
174     }
175 
176     private static class FactoryData extends AbstractFactoryData {
177         private final InetSocketAddress[] contactPoints;
178         private final ColumnMapping[] columns;
179         private final boolean useTls;
180         private final String clusterName;
181         private final String keyspace;
182         private final String table;
183         private final String username;
184         private final String password;
185         private final boolean useClockForTimestampGenerator;
186         private final boolean batched;
187         private final BatchStatement.Type batchType;
188 
189         private FactoryData(final SocketAddress[] contactPoints, final ColumnMapping[] columns, final boolean useTls,
190                             final String clusterName, final String keyspace, final String table, final String username,
191                             final String password, final boolean useClockForTimestampGenerator, final int bufferSize,
192                             final boolean batched, final BatchStatement.Type batchType) {
193             super(bufferSize);
194             this.contactPoints = convertAndAddDefaultPorts(contactPoints);
195             this.columns = columns;
196             this.useTls = useTls;
197             this.clusterName = clusterName;
198             this.keyspace = keyspace;
199             this.table = table;
200             this.username = username;
201             this.password = password;
202             this.useClockForTimestampGenerator = useClockForTimestampGenerator;
203             this.batched = batched;
204             this.batchType = batchType;
205         }
206 
207         private static InetSocketAddress[] convertAndAddDefaultPorts(final SocketAddress... socketAddresses) {
208             final InetSocketAddress[] inetSocketAddresses = new InetSocketAddress[socketAddresses.length];
209             for (int i = 0; i < inetSocketAddresses.length; i++) {
210                 final SocketAddress socketAddress = socketAddresses[i];
211                 inetSocketAddresses[i] = socketAddress.getPort() == 0
212                     ? new InetSocketAddress(socketAddress.getAddress(), DEFAULT_PORT)
213                     : socketAddress.getSocketAddress();
214             }
215             return inetSocketAddresses;
216         }
217     }
218 }