This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git
The following commit(s) were added to refs/heads/master by this push:
new b3a904e [INLONG-2054][feature][audit] audit-sdk add disaster recovery
(#2198)
b3a904e is described below
commit b3a904e9db8bc8bee0ae6be7919f1d37e98b6b6d
Author: doleyzi <[email protected]>
AuthorDate: Thu Jan 20 12:04:00 2022 +0800
[INLONG-2054][feature][audit] audit-sdk add disaster recovery (#2198)
---
.../java/org/apache/inlong/audit/AuditImp.java | 21 +++-
.../org/apache/inlong/audit/send/SenderGroup.java | 2 +-
.../apache/inlong/audit/send/SenderManager.java | 137 +++++++++++++++++++--
.../org/apache/inlong/audit/util/AuditConfig.java | 81 ++++++++++++
.../org/apache/inlong/audit/util/AuditData.java | 3 +-
.../apache/inlong/audit/send/SenderGroupTest.java | 4 +-
.../inlong/audit/send/SenderManagerTest.java | 5 +-
7 files changed, 230 insertions(+), 23 deletions(-)
diff --git
a/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/AuditImp.java
b/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/AuditImp.java
index 872cc78..28b1b8b 100644
--- a/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/AuditImp.java
+++ b/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/AuditImp.java
@@ -19,6 +19,7 @@ package org.apache.inlong.audit;
import org.apache.inlong.audit.protocol.AuditApi;
import org.apache.inlong.audit.send.SenderManager;
+import org.apache.inlong.audit.util.AuditConfig;
import org.apache.inlong.audit.util.Config;
import org.apache.inlong.audit.util.StatInfo;
import org.slf4j.Logger;
@@ -45,6 +46,7 @@ public class AuditImp {
private HashMap<String, StatInfo> threadSumMap = new HashMap<String,
StatInfo>();
private ConcurrentHashMap<String, StatInfo> deleteCountMap = new
ConcurrentHashMap<String, StatInfo>();
private List<String> deleteKeyList = new ArrayList<String>();
+ private AuditConfig auditConfig = null;
private Config config = new Config();
private Long sdkTime;
private int packageId = 1;
@@ -79,7 +81,10 @@ public class AuditImp {
}
config.init();
timer.schedule(timerTask, PERIOD, PERIOD);
- this.manager = new SenderManager(config);
+ if (auditConfig == null) {
+ auditConfig = new AuditConfig();
+ }
+ this.manager = new SenderManager(auditConfig);
}
/**
@@ -103,6 +108,16 @@ public class AuditImp {
}
/**
+ * set audit config
+ *
+ * @param config
+ */
+ public void setAuditConfig(AuditConfig config) {
+ auditConfig = config;
+ manager.setAuditConfig(config);
+ }
+
+ /**
* api
*
* @param auditID
@@ -178,7 +193,7 @@ public class AuditImp {
requestBulid.setMsgHeader(mssageHeader).setRequestId(manager.nextRequestId());
for (Map.Entry<String, StatInfo> entry : threadSumMap.entrySet()) {
String[] keyArray = entry.getKey().split(FIELD_SEPARATORS);
- long logTime = Long.parseLong(keyArray[0]) * 60;
+ long logTime = Long.parseLong(keyArray[0]) * PERIOD;
String inlongGroupID = keyArray[1];
String inlongStreamID = keyArray[2];
String auditID = keyArray[3];
@@ -202,7 +217,7 @@ public class AuditImp {
requestBulid.clearMsgBody();
}
threadSumMap.clear();
- logger.info("finished send report.");
+ logger.info("finish send report.");
}
/**
diff --git
a/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/send/SenderGroup.java
b/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/send/SenderGroup.java
index fab58e8..c578f06 100644
---
a/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/send/SenderGroup.java
+++
b/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/send/SenderGroup.java
@@ -132,7 +132,7 @@ public class SenderGroup {
try {
Thread.sleep(waitChannelIntervalMs);
} catch (Throwable e) {
- System.out.println(e.getMessage());
+ logger.error(e.getMessage());
}
}
if (channel == null) {
diff --git
a/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/send/SenderManager.java
b/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/send/SenderManager.java
index 6959341..ff37165 100644
---
a/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/send/SenderManager.java
+++
b/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/send/SenderManager.java
@@ -18,8 +18,8 @@
package org.apache.inlong.audit.send;
import org.apache.inlong.audit.protocol.AuditApi;
+import org.apache.inlong.audit.util.AuditConfig;
import org.apache.inlong.audit.util.AuditData;
-import org.apache.inlong.audit.util.Config;
import org.apache.inlong.audit.util.Decoder;
import org.apache.inlong.audit.util.SenderResult;
import org.jboss.netty.buffer.ChannelBuffer;
@@ -30,10 +30,17 @@ import org.jboss.netty.channel.MessageEvent;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.io.File;
+import java.io.FileInputStream;
+import java.io.FileOutputStream;
+import java.io.IOException;
+import java.io.ObjectInputStream;
+import java.io.ObjectOutputStream;
import java.security.SecureRandom;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
+import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
@@ -44,9 +51,10 @@ public class SenderManager {
private static final Logger logger =
LoggerFactory.getLogger(SenderManager.class);
public static final int DEFAULT_SEND_THREADNUM = 2;
public static final Long MAX_REQUEST_ID = 1000000000L;
+ private static final int SEND_INTERVAL_MS = 20;
public static final int ALL_CONNECT_CHANNEL = -1;
public static final int DEFAULT_CONNECT_CHANNEL = 2;
-
+ public static final String LF = "\n";
private SenderGroup sender;
private int maxConnectChannels = ALL_CONNECT_CHANNEL;
private SecureRandom sRandom = new
SecureRandom(Long.toString(System.currentTimeMillis()).getBytes());
@@ -54,14 +62,14 @@ public class SenderManager {
private HashSet<String> currentIpPorts = new HashSet<String>();
private AtomicLong requestIdSeq = new AtomicLong(0L);
private ConcurrentHashMap<Long, AuditData> dataMap = new
ConcurrentHashMap<>();
- private Config config;
+ private AuditConfig auditConfig;
/**
* Constructor
*
* @param config
*/
- public SenderManager(Config config) {
+ public SenderManager(AuditConfig config) {
this(config, DEFAULT_CONNECT_CHANNEL);
}
@@ -71,9 +79,9 @@ public class SenderManager {
* @param config
* @param maxConnectChannels
*/
- public SenderManager(Config config, int maxConnectChannels) {
+ public SenderManager(AuditConfig config, int maxConnectChannels) {
try {
- this.config = config;
+ this.auditConfig = config;
this.maxConnectChannels = maxConnectChannels;
SenderHandler clientHandler = new SenderHandler(this);
this.sender = new SenderGroup(DEFAULT_SEND_THREADNUM, new
Decoder(), clientHandler);
@@ -136,7 +144,7 @@ public class SenderManager {
AuditData data = new AuditData(sdkTime, baseCommand);
// Cache first
this.dataMap.putIfAbsent(baseCommand.getAuditRequest().getRequestId(),
data);
- this.sendData(data);
+ this.sendData(data.getDataByte());
}
/**
@@ -144,8 +152,8 @@ public class SenderManager {
*
* @param data
*/
- private void sendData(AuditData data) {
- ChannelBuffer dataBuf =
ChannelBuffers.wrappedBuffer(data.getDataByte());
+ private void sendData(byte[] data) {
+ ChannelBuffer dataBuf = ChannelBuffers.wrappedBuffer(data);
SenderResult result = this.sender.send(dataBuf);
if (!result.result) {
this.sender.setHasSendError(true);
@@ -156,8 +164,92 @@ public class SenderManager {
* Clean up the backlog of unsent message packets
*/
public void clearBuffer() {
+ logger.info("failed cache size:" + this.dataMap.size());
for (AuditData data : this.dataMap.values()) {
- this.sendData(data);
+ this.sendData(data.getDataByte());
+ sleep(SEND_INTERVAL_MS);
+ }
+ if (this.dataMap.size() == 0) {
+ checkAuditFile();
+ }
+ if (this.dataMap.size() > auditConfig.getMaxCacheRow()) {
+ logger.info("failed cache size: {}>{}", this.dataMap.size(),
auditConfig.getMaxCacheRow());
+ writeLocalFile();
+ this.dataMap.clear();
+ }
+ }
+
+ /**
+ * write local file
+ */
+ private void writeLocalFile() {
+ try {
+ if (!checkFilePath()) {
+ return;
+ }
+ File file = new File(auditConfig.getDisasterFile());
+ if (!file.exists()) {
+ if (!file.createNewFile()) {
+ logger.error("create {} {}",
auditConfig.getDisasterFile(), " failed");
+ return;
+ }
+ logger.info("create {}", auditConfig.getDisasterFile());
+ }
+ if (file.length() > auditConfig.getMaxFileSize()) {
+ file.delete();
+ return;
+ }
+ FileOutputStream fos = new FileOutputStream(file);
+ ObjectOutputStream objectOutputStream = new
ObjectOutputStream(fos);
+ objectOutputStream.writeObject(dataMap);
+ objectOutputStream.close();
+ fos.close();
+ } catch (IOException ioException) {
+ logger.error(ioException.getMessage());
+ }
+ }
+
+ /**
+ * check file path
+ *
+ * @return
+ */
+ private boolean checkFilePath() {
+ File file = new File(auditConfig.getFilePath());
+ if (!file.exists()) {
+ if (!file.mkdirs()) {
+ return false;
+ }
+ logger.info("create {}", auditConfig.getFilePath());
+ }
+ return true;
+ }
+
+ /**
+ * check audit file
+ */
+ private void checkAuditFile() {
+ try {
+ File file = new File(auditConfig.getDisasterFile());
+ if (!file.exists()) {
+ return;
+ }
+ FileInputStream inputStream = new
FileInputStream(auditConfig.getDisasterFile());
+ ObjectInputStream objectInputStream = new
ObjectInputStream(inputStream);
+ ConcurrentHashMap<Long, AuditData> fileData =
+ (ConcurrentHashMap<Long, AuditData>)
objectInputStream.readObject();
+ for (Map.Entry<Long, AuditData> entry : fileData.entrySet()) {
+ if (this.dataMap.size() < (auditConfig.getMaxCacheRow() / 2)) {
+ this.dataMap.putIfAbsent(entry.getKey(), entry.getValue());
+ }
+ this.sendData(entry.getValue().getDataByte());
+ sleep(SEND_INTERVAL_MS);
+ }
+ objectInputStream.close();
+ inputStream.close();
+ file.delete();
+ } catch (IOException | ClassNotFoundException ioException) {
+ logger.error(ioException.getMessage());
}
}
@@ -193,14 +285,12 @@ public class SenderManager {
}
if
(AuditApi.AuditReply.RSP_CODE.SUCCESS.equals(baseCommand.getAuditReply().getRspCode()))
{
this.dataMap.remove(requestId);
- this.sender.notifyAll();
return;
}
int resendTimes = data.increaseResendTimes();
if (resendTimes <
org.apache.inlong.audit.send.SenderGroup.MAX_SEND_TIMES) {
- this.sendData(data);
+ this.sendData(data.getDataByte());
}
- this.sender.notifyAll();
} catch (Throwable ex) {
logger.error(ex.getMessage());
this.sender.setHasSendError(true);
@@ -221,4 +311,25 @@ public class SenderManager {
logger.error(ex.getMessage());
}
}
+
+ /**
+ * sleep
+ *
+ * @param millisecond
+ */
+ private void sleep(int millisecond) {
+ try {
+ Thread.sleep(millisecond);
+ } catch (Throwable e) {
+ logger.error(e.getMessage());
+ }
+ }
+
+ /***
+ * set audit config
+ * @param config
+ */
+ public void setAuditConfig(AuditConfig config) {
+ auditConfig = config;
+ }
}
diff --git
a/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/util/AuditConfig.java
b/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/util/AuditConfig.java
new file mode 100644
index 0000000..67e8f94
--- /dev/null
+++
b/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/util/AuditConfig.java
@@ -0,0 +1,81 @@
+/*
+ * 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.inlong.audit.util;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+public class AuditConfig {
+ private static final Logger logger =
LoggerFactory.getLogger(AuditConfig.class);
+ private static String FILE_PATH = "/data/inlong/audit/";
+ private static final int FILE_SIZE = 500 * 1024 * 1024;
+ private static final int MAX_CACHE_ROWS = 2000000;
+ private static final int MIN_CACHE_ROWS = 100;
+
+ private String filePath;
+ private int maxCacheRow;
+ private int maxFileSize = FILE_SIZE;
+
+ public AuditConfig(String filePath, int maxCacheRow) {
+ if (filePath == null || filePath.length() == 0) {
+ this.filePath = FILE_PATH;
+ } else {
+ this.filePath = filePath;
+ }
+ if (maxCacheRow < MIN_CACHE_ROWS) {
+ this.maxCacheRow = MAX_CACHE_ROWS;
+ } else {
+ this.maxCacheRow = maxCacheRow;
+ }
+ }
+
+ public AuditConfig() {
+ this.filePath = FILE_PATH;
+ this.maxCacheRow = MAX_CACHE_ROWS;
+ }
+
+ public void setFilePath(String filePath) {
+ this.filePath = filePath;
+ }
+
+ public void setMaxCacheRow(int maxCacheRow) {
+ this.maxCacheRow = maxCacheRow;
+ }
+
+ public String getFilePath() {
+ return filePath;
+ }
+
+ public int getMaxCacheRow() {
+ return maxCacheRow;
+ }
+
+ public int getMaxFileSize() {
+ return maxFileSize;
+ }
+
+ public void setMaxFileSize(int maxFileSize) {
+ this.maxFileSize = maxFileSize;
+ }
+
+ public String getDisasterFile() {
+ return filePath + "/" + disasterFileName;
+ }
+
+ private String disasterFileName = "disaster.data";
+}
diff --git
a/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/util/AuditData.java
b/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/util/AuditData.java
index 03c101e..4dde7e1 100644
---
a/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/util/AuditData.java
+++
b/inlong-audit/audit-sdk/src/main/java/org/apache/inlong/audit/util/AuditData.java
@@ -19,10 +19,11 @@ package org.apache.inlong.audit.util;
import org.apache.inlong.audit.protocol.AuditApi;
+import java.io.Serializable;
import java.nio.ByteBuffer;
import java.util.concurrent.atomic.AtomicInteger;
-public class AuditData {
+public class AuditData implements Serializable {
public static int HEAD_LENGTH = 4;
private final long sdkTime;
private final AuditApi.BaseCommand content;
diff --git
a/inlong-audit/audit-sdk/src/test/java/org/apache/inlong/audit/send/SenderGroupTest.java
b/inlong-audit/audit-sdk/src/test/java/org/apache/inlong/audit/send/SenderGroupTest.java
index 2b1a8f0..4cc8192 100644
---
a/inlong-audit/audit-sdk/src/test/java/org/apache/inlong/audit/send/SenderGroupTest.java
+++
b/inlong-audit/audit-sdk/src/test/java/org/apache/inlong/audit/send/SenderGroupTest.java
@@ -17,7 +17,7 @@
package org.apache.inlong.audit.send;
-import org.apache.inlong.audit.util.Config;
+import org.apache.inlong.audit.util.AuditConfig;
import org.apache.inlong.audit.util.Decoder;
import org.junit.Test;
@@ -25,7 +25,7 @@ import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
public class SenderGroupTest {
- Config testConfig = new Config();
+ AuditConfig testConfig = new AuditConfig();
SenderManager testManager = new SenderManager(testConfig);
SenderHandler clientHandler = new
org.apache.inlong.audit.send.SenderHandler(testManager);
SenderGroup sender = new org.apache.inlong.audit.send.SenderGroup(10, new
Decoder(), clientHandler);
diff --git
a/inlong-audit/audit-sdk/src/test/java/org/apache/inlong/audit/send/SenderManagerTest.java
b/inlong-audit/audit-sdk/src/test/java/org/apache/inlong/audit/send/SenderManagerTest.java
index 6466149..4ddf3d6 100644
---
a/inlong-audit/audit-sdk/src/test/java/org/apache/inlong/audit/send/SenderManagerTest.java
+++
b/inlong-audit/audit-sdk/src/test/java/org/apache/inlong/audit/send/SenderManagerTest.java
@@ -17,19 +17,18 @@
package org.apache.inlong.audit.send;
-import org.apache.inlong.audit.util.Config;
+import org.apache.inlong.audit.util.AuditConfig;
import org.junit.Test;
import static org.junit.Assert.assertTrue;
public class SenderManagerTest {
- private Config testConfig = new Config();
+ private AuditConfig testConfig = new AuditConfig();
@Test
public void nextRequestId() {
SenderManager testManager = new SenderManager(testConfig);
Long requestId = testManager.nextRequestId();
- System.out.println(requestId);
assertTrue(requestId == 0);
requestId = testManager.nextRequestId();