1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
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
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
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
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 }