Java 写入 influxdb
·
利用Python随机生成一个1000行的csv文件
import csv
import random
from datetime import datetime, timedelta
from random import randint, choice
# 定义监控对象列表和指标名称列表
monitor_objects = ['Server1', 'Server2', 'Server3', 'DB1']
metric_names = ['CPUUsage', 'MemoryUsage', 'DiskSpace', 'NetworkTraffic']
# 生成随机日期范围,例如从今天开始往回计算一年的数据
start_date = datetime.now() - timedelta(days=365)
end_date = datetime.now()
# 打开一个文件,准备写入
filename = 'random_data.csv'
with open(filename, mode='w', newline='', encoding='utf-8') as csvfile:
fieldnames = ['dates', 'monitorObject', 'timeStamp', 'metricName', 'value']
writer = csv.DictWriter(csvfile, fieldnames=fieldnames)
# 写入表头
writer.writeheader()
# 生成并写入1000行数据
for _ in range(10000000):
# 随机日期
random_date = start_date + (end_date - start_date) * random.random()
date_str = random_date.strftime('%Y-%m-%d')
time_stamp = int(random_date.timestamp())
# 随机选择监控对象和指标名称
monitor_object = choice(monitor_objects)
metric_name = choice(metric_names)
# 随机值,根据实际情况调整范围
value = round(random.uniform(0, 100), 2) # 假设值在0到100之间,保留两位小数
# 写入一行数据
writer.writerow({
'dates': date_str,
'monitorObject': monitor_object,
'timeStamp': time_stamp,
'metricName': metric_name,
'value': value
})
print(f"CSV文件已成功生成,文件名为: {filename}")
利用Java写入inluxdb
package org.example;
import com.csvreader.CsvReader;
import com.influxdb.client.InfluxDBClient;
import com.influxdb.client.InfluxDBClientFactory;
import com.influxdb.client.WriteApiBlocking;
import com.influxdb.client.domain.WritePrecision;
import com.influxdb.client.write.Point;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.List;
public class Main {
static String token = "ZgpSCp3H_9liDB6v5POacJ8MOqnKOUR9YUlJjkIvLtFbYVyZr5Rkn-YKpUNdPqWpRKY_7Fqdwv6GN9r4A8BwcQ==";
static String bucket = "data";
static String org = "test";
static InfluxDBClient client = InfluxDBClientFactory.create("http://localhost:8086", token.toCharArray());
public static void main(String[] args) throws Exception {
String CSV_FILE_PATH = "D:\\JAVA_dataAnalyze\\dataAnalyze\\src\\main\\java\\org\\example\\data\\random_data.csv";
long startTime = System.currentTimeMillis();
readCsv(CSV_FILE_PATH); // 方法1 串行写入
long endTime = System.currentTimeMillis();
System.out.println("excute time: " + (endTime - startTime) + "ms");
startTime = System.currentTimeMillis();
readCSVByMemory(CSV_FILE_PATH); // 方法2 stream流并行写入
endTime = System.currentTimeMillis();
System.out.println("excute time: " + (endTime - startTime) + "ms");
}
public static void readCsv(String CSV_FILE_PATH) throws Exception {
CsvReader CSVReader = new CsvReader(CSV_FILE_PATH, ',', StandardCharsets.UTF_8);
CSVReader.readHeaders();
int rowNumber = 0;
long startTime = System.currentTimeMillis();
while (CSVReader.readRecord()) {
// 安装mysql并引入Mapping
String content = CSVReader.getRawRecord();
// 按照 dates,monitorObject,timeStamp,metricName,value 解析内容
Point lineProtocol = row2InfluxPoint(content);
//System.out.println(lineProtocol.toLineProtocol());
point2InluxDB(lineProtocol);
rowNumber++;
if(rowNumber % 10000 == 0){
long endTime = System.currentTimeMillis();
System.out.println("rowNumber: " + rowNumber + "excute time: " + (endTime - startTime) + "ms");
startTime = endTime;
}
}
CSVReader.close();
}
public static void readCSVByMemory(String CSV_FILE_PATH) throws Exception {
CsvReader CSVReader = new CsvReader(CSV_FILE_PATH, ',', StandardCharsets.UTF_8);
CSVReader.readHeaders();
int rowNumber = 0;
long startTime = System.currentTimeMillis();
List<String> contentList = new ArrayList<String>();
while (CSVReader.readRecord()) {
contentList.add(CSVReader.getRawRecord());
}
// 并行将contentList中的元素写入influxdb
contentList.stream().parallel().forEach(content -> {
point2InluxDB(row2InfluxPoint(content));
});
}
public static Point row2InfluxPoint(String content) {
String[] contentArray = content.split(",");
Point point = Point
.measurement("mem_data1T")
.addTag("monitorObject", contentArray[1])
.addField(contentArray[3], Float.parseFloat(contentArray[4]))
.time(Long.parseLong(contentArray[2]), WritePrecision.S);
return point;
}
public static void point2InluxDB(Point point){
WriteApiBlocking writeApi = client.getWriteApiBlocking();
writeApi.writePoint(bucket, org, point);
}
}
<dependencies>
<dependency>
<groupId>net.sourceforge.javacsv</groupId>
<artifactId>javacsv</artifactId>
<version>2.0</version>
</dependency>
<dependency>
<groupId>com.influxdb</groupId>
<artifactId>influxdb-client-java</artifactId>
<version>6.6.0</version>
</dependency>
</dependencies>
更多推荐
所有评论(0)