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();

Reply via email to