mirror of
https://github.com/abcv7/sfbx-cloud.git
synced 2026-08-16 11:47:00 +00:00
first commit
This commit is contained in:
@@ -0,0 +1,41 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project xmlns="http://maven.apache.org/POM/4.0.0"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
|
||||
<parent>
|
||||
<artifactId>sfbx-framework</artifactId>
|
||||
<groupId>com.itheima.sfbx</groupId>
|
||||
<version>2.0-SNAPSHOT</version>
|
||||
</parent>
|
||||
<modelVersion>4.0.0</modelVersion>
|
||||
|
||||
<artifactId>framework-influxdb</artifactId>
|
||||
|
||||
<dependencies>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-web</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-autoconfigure</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>com.itheima.sfbx</groupId>
|
||||
<artifactId>framework-commons</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-configuration-processor</artifactId>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.influxdb</groupId>
|
||||
<artifactId>influxdb-java</artifactId>
|
||||
<version>2.18</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>org.aspectj</groupId>
|
||||
<artifactId>aspectjweaver</artifactId>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
</project>
|
||||
+14
@@ -0,0 +1,14 @@
|
||||
package com.itheima.sfbx.framework.influxdb;
|
||||
|
||||
import com.itheima.sfbx.framework.influxdb.anno.Insert;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
public interface InfluxDBBaseMapper<T> {
|
||||
|
||||
@Insert
|
||||
void insertOne(T entity);
|
||||
|
||||
@Insert
|
||||
void insertBatch(List<T> entityList);
|
||||
}
|
||||
+11
@@ -0,0 +1,11 @@
|
||||
package com.itheima.sfbx.framework.influxdb.anno;
|
||||
|
||||
import java.lang.annotation.ElementType;
|
||||
import java.lang.annotation.Retention;
|
||||
import java.lang.annotation.RetentionPolicy;
|
||||
import java.lang.annotation.Target;
|
||||
|
||||
@Target(ElementType.METHOD)
|
||||
@Retention(RetentionPolicy.RUNTIME)
|
||||
public @interface Insert {
|
||||
}
|
||||
+10
@@ -0,0 +1,10 @@
|
||||
package com.itheima.sfbx.framework.influxdb.anno;
|
||||
|
||||
import java.lang.annotation.*;
|
||||
|
||||
@Retention(RetentionPolicy.RUNTIME)
|
||||
@Target({ElementType.PARAMETER})
|
||||
@Documented
|
||||
public @interface Param {
|
||||
String value();
|
||||
}
|
||||
+19
@@ -0,0 +1,19 @@
|
||||
package com.itheima.sfbx.framework.influxdb.anno;
|
||||
|
||||
import java.lang.annotation.ElementType;
|
||||
import java.lang.annotation.Retention;
|
||||
import java.lang.annotation.RetentionPolicy;
|
||||
import java.lang.annotation.Target;
|
||||
|
||||
@Target(ElementType.METHOD)
|
||||
@Retention(RetentionPolicy.RUNTIME)
|
||||
public @interface Select {
|
||||
//执行的influxQL
|
||||
String value();
|
||||
|
||||
//返回的类型
|
||||
Class resultType();
|
||||
|
||||
//执行的目标库
|
||||
String database();
|
||||
}
|
||||
+68
@@ -0,0 +1,68 @@
|
||||
package com.itheima.sfbx.framework.influxdb.aspect;
|
||||
|
||||
import com.itheima.sfbx.framework.influxdb.anno.Insert;
|
||||
import com.itheima.sfbx.framework.influxdb.anno.Select;
|
||||
import com.itheima.sfbx.framework.influxdb.core.Executor;
|
||||
import com.itheima.sfbx.framework.influxdb.core.ParameterHandler;
|
||||
import com.itheima.sfbx.framework.influxdb.core.ResultSetHandler;
|
||||
import org.aspectj.lang.ProceedingJoinPoint;
|
||||
import org.aspectj.lang.annotation.Around;
|
||||
import org.aspectj.lang.annotation.Aspect;
|
||||
import org.aspectj.lang.reflect.MethodSignature;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.lang.reflect.Method;
|
||||
import java.lang.reflect.Parameter;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* @ClassName InfluxDBAspect.java
|
||||
* @Description TODO
|
||||
*/
|
||||
@Aspect
|
||||
@Component
|
||||
public class InfluxDBAspect {
|
||||
|
||||
private final Executor executor;
|
||||
|
||||
private final ParameterHandler parameterHandler;
|
||||
|
||||
private final ResultSetHandler resultSetHandler;
|
||||
|
||||
@Autowired
|
||||
public InfluxDBAspect(Executor executor,ParameterHandler parameterHandler,ResultSetHandler resultSetHandler) {
|
||||
this.executor = executor;
|
||||
this.parameterHandler = parameterHandler;
|
||||
this.resultSetHandler = resultSetHandler;
|
||||
}
|
||||
|
||||
@Around("@annotation(select)")
|
||||
public Object select(ProceedingJoinPoint joinPoint, Select select) {
|
||||
MethodSignature methodSignature = (MethodSignature) joinPoint.getSignature();
|
||||
Method method = methodSignature.getMethod();
|
||||
Select selectAnnotation = method.getAnnotation(Select.class);
|
||||
//获得执行参数
|
||||
Parameter[] parameters = method.getParameters();
|
||||
//获得执行参数值
|
||||
Object[] args = joinPoint.getArgs();
|
||||
//获得执行sql
|
||||
String sql = selectAnnotation.value();
|
||||
//替换参数
|
||||
sql = parameterHandler.handleParameter(parameters,args,sql);
|
||||
//注解声明返回类型
|
||||
Class<?> resultType = selectAnnotation.resultType();
|
||||
//查询结果
|
||||
List<Map<String,Object>> reultList = executor.select(sql,selectAnnotation.database());
|
||||
//根据返回类型返回结果
|
||||
return resultSetHandler.handleResultSet(reultList, method,sql,resultType);
|
||||
}
|
||||
|
||||
@Around("@annotation(insert)")
|
||||
public void insert(ProceedingJoinPoint joinPoint, Insert insert) {
|
||||
//获得执行参数值
|
||||
Object[] args = joinPoint.getArgs();
|
||||
executor.insert(args);
|
||||
}
|
||||
}
|
||||
+31
@@ -0,0 +1,31 @@
|
||||
package com.itheima.sfbx.framework.influxdb.config;
|
||||
|
||||
|
||||
import com.itheima.sfbx.framework.influxdb.core.Executor;
|
||||
import com.itheima.sfbx.framework.influxdb.core.ParameterHandler;
|
||||
import com.itheima.sfbx.framework.influxdb.core.ResultSetHandler;
|
||||
import org.influxdb.InfluxDB;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
|
||||
/**
|
||||
* 时序数据库配置类
|
||||
*/
|
||||
@Configuration
|
||||
public class InfluxDBConfig {
|
||||
|
||||
@Bean(name = "executor")
|
||||
public Executor executor(InfluxDB influxDB) {
|
||||
return new Executor(influxDB);
|
||||
}
|
||||
|
||||
@Bean(name = "parameterHandler")
|
||||
public ParameterHandler parameterHandler(InfluxDB influxDB) {
|
||||
return new ParameterHandler();
|
||||
}
|
||||
|
||||
@Bean(name = "resultSetHandler")
|
||||
public ResultSetHandler resultSetHandler(InfluxDB influxDB) {
|
||||
return new ResultSetHandler();
|
||||
}
|
||||
}
|
||||
+114
@@ -0,0 +1,114 @@
|
||||
package com.itheima.sfbx.framework.influxdb.core;
|
||||
|
||||
import com.google.common.collect.Lists;
|
||||
import com.itheima.sfbx.framework.commons.utils.EmptyUtil;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.influxdb.InfluxDB;
|
||||
import org.influxdb.annotation.Measurement;
|
||||
import org.influxdb.dto.BatchPoints;
|
||||
import org.influxdb.dto.Point;
|
||||
import org.influxdb.dto.Query;
|
||||
import org.influxdb.dto.QueryResult;
|
||||
import org.influxdb.impl.InfluxDBMapper;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* 执行器
|
||||
*/
|
||||
@Slf4j
|
||||
public class Executor {
|
||||
|
||||
InfluxDB influxDB;
|
||||
|
||||
public Executor() {
|
||||
}
|
||||
|
||||
public Executor(InfluxDB influxDB) {
|
||||
this.influxDB = influxDB;
|
||||
}
|
||||
|
||||
public List<Map<String,Object>> select(String sql,String database) {
|
||||
QueryResult queryResult = influxDB.query(new Query(sql, database));
|
||||
List<Map<String,Object>> resultList = new ArrayList<>();
|
||||
queryResult.getResults().forEach(result -> {
|
||||
//查询出错抛出错误信息
|
||||
if (!EmptyUtil.isNullOrEmpty(result.getError())){
|
||||
throw new RuntimeException(result.getError());
|
||||
}
|
||||
if (!EmptyUtil.isNullOrEmpty(result)&&!EmptyUtil.isNullOrEmpty(result.getSeries())){
|
||||
//获取所有列的集合,一个迭代是代表一组
|
||||
List<QueryResult.Series> series= result.getSeries();
|
||||
for (QueryResult.Series s : series) {
|
||||
//列中含有多行数据,每行数据含有多列value,所以嵌套List
|
||||
List<List<Object>> values = s.getValues();
|
||||
//每组的列是固定的
|
||||
List<String> columns = s.getColumns();
|
||||
for (List<Object> v:values){
|
||||
//循环遍历结果集,获取每行对应的value,以map形式保存
|
||||
Map<String,Object> queryMap =new HashMap<String, Object>();
|
||||
for(int i=0;i<columns.size();i++){
|
||||
//遍历所有列名,获取列对应的值
|
||||
String column = columns.get(i);
|
||||
if (v.get(i)==null||v.get(i).equals("null")){
|
||||
//如果是null就存入null
|
||||
queryMap.put(column,null);
|
||||
}else {
|
||||
//不是null就转成字符串存储
|
||||
String value = String.valueOf(v.get(i));
|
||||
//如果是时间戳还可以格式转换,我这里懒了
|
||||
queryMap.put(column, value);
|
||||
}
|
||||
}
|
||||
//把结果添加到结果集中
|
||||
resultList.add(queryMap);
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
return resultList;
|
||||
}
|
||||
|
||||
public void insert(Object args[]) {
|
||||
if (args.length != 1) {
|
||||
throw new RuntimeException();
|
||||
}
|
||||
Object obj = args[0];
|
||||
List<Object> list = Lists.newArrayList();
|
||||
if (obj instanceof List){
|
||||
list = (ArrayList) obj;
|
||||
}else {
|
||||
list.add(obj);
|
||||
}
|
||||
if (list.size() > 0) {
|
||||
Object firstObj = list.get(0);
|
||||
Class<?> domainClass = firstObj.getClass();
|
||||
List<Point> pointList = new ArrayList<>();
|
||||
for (Object o : list) {
|
||||
Point point = Point
|
||||
.measurementByPOJO(domainClass)
|
||||
.addFieldsFromPOJO(o)
|
||||
.build();
|
||||
pointList.add(point);
|
||||
}
|
||||
//获取数据库名和rp
|
||||
Measurement measurement = firstObj.getClass().getAnnotation(Measurement.class);
|
||||
String database = measurement.database();
|
||||
String retentionPolicy = measurement.retentionPolicy();
|
||||
BatchPoints batchPoints = BatchPoints
|
||||
.builder()
|
||||
.points(pointList)
|
||||
.retentionPolicy(retentionPolicy).build();
|
||||
influxDB.setDatabase(database);
|
||||
influxDB.write(batchPoints);
|
||||
}
|
||||
}
|
||||
|
||||
public void delete(String sql, String database) {
|
||||
influxDB.query(new Query(sql, database));
|
||||
}
|
||||
|
||||
}
|
||||
+40
@@ -0,0 +1,40 @@
|
||||
package com.itheima.sfbx.framework.influxdb.core;
|
||||
|
||||
import com.itheima.sfbx.framework.influxdb.anno.Param;
|
||||
|
||||
import java.lang.reflect.Parameter;
|
||||
|
||||
/**
|
||||
* 参数处理器
|
||||
*/
|
||||
public class ParameterHandler {
|
||||
|
||||
/**
|
||||
* 拼接sql
|
||||
*
|
||||
* @param parameters 参数名
|
||||
* @param args 参数实际值
|
||||
* @param sql 未拼接参数的sql语句
|
||||
* @return 拼接好的sql
|
||||
*/
|
||||
public String handleParameter(Parameter[] parameters, Object[] args, String sql) {
|
||||
for (int i = 0; i < parameters.length; i++) {
|
||||
Class<?> parameterType = parameters[i].getType();
|
||||
String parameterName = parameters[i].getName();
|
||||
|
||||
Param param = parameters[i].getAnnotation(Param.class);
|
||||
if (param != null) {
|
||||
parameterName = param.value();
|
||||
}
|
||||
|
||||
if (parameterType == String.class) {
|
||||
sql = sql.replaceAll("\\#\\{" + parameterName + "\\}", "'" + args[i] + "'");
|
||||
sql = sql.replaceAll("\\$\\{" + parameterName + "\\}", args[i].toString());
|
||||
} else {
|
||||
sql = sql.replaceAll("\\#\\{" + parameterName + "\\}", args[i].toString());
|
||||
sql = sql.replaceAll("\\$\\{" + parameterName + "\\}", args[i].toString());
|
||||
}
|
||||
}
|
||||
return sql;
|
||||
}
|
||||
}
|
||||
+183
@@ -0,0 +1,183 @@
|
||||
package com.itheima.sfbx.framework.influxdb.core;
|
||||
|
||||
import com.google.common.collect.Lists;
|
||||
import com.google.common.collect.Maps;
|
||||
import com.itheima.sfbx.framework.commons.utils.BeanConv;
|
||||
import com.itheima.sfbx.framework.commons.utils.EmptyUtil;
|
||||
import lombok.SneakyThrows;
|
||||
|
||||
import java.lang.reflect.Constructor;
|
||||
import java.lang.reflect.InvocationTargetException;
|
||||
import java.lang.reflect.Method;
|
||||
import java.math.BigDecimal;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* 结果集处理器
|
||||
*/
|
||||
public class ResultSetHandler {
|
||||
|
||||
/***
|
||||
* @description 结果处理
|
||||
*
|
||||
* @param reultList influx返回结果
|
||||
* @param method 目标方法
|
||||
* @param sql 执行sql
|
||||
* @param resultType 注解声明返回类型
|
||||
* @return
|
||||
* @return: java.lang.Object
|
||||
*/
|
||||
@SneakyThrows
|
||||
public Object handleResultSet(List<Map<String,Object>> reultList, Method method, String sql, Class resultType) {
|
||||
Class<?> returnTypeTarget = method.getReturnType();
|
||||
//如果结果为空直接返回空构建
|
||||
if (EmptyUtil.isNullOrEmpty(reultList)){
|
||||
if (returnTypeTarget== List.class){
|
||||
return Lists.newArrayList();
|
||||
}else if (returnTypeTarget==Map.class){
|
||||
return Maps.newHashMap();
|
||||
}else if (returnTypeTarget==String.class){
|
||||
return null;
|
||||
}else {
|
||||
return convertStringToObject(resultType,"0");
|
||||
}
|
||||
}
|
||||
//当前method声明返回结果不为list,且resultType与method声明返回结果类型不匹配
|
||||
if (returnTypeTarget!= List.class&&resultType!=returnTypeTarget){
|
||||
throw new RuntimeException("返回类型与声明返回类型不匹配");
|
||||
}
|
||||
//当前method声明返回结果不为list,且resultType与method声明返回结果类型匹配
|
||||
if (returnTypeTarget!= List.class&&resultType==returnTypeTarget){
|
||||
//结果不唯一则抛出异常
|
||||
if (reultList.size()!=1){
|
||||
throw new RuntimeException("返回结果不唯一");
|
||||
}
|
||||
//驼峰处理
|
||||
Map<String, Object> mapHandler = convertKeysToCamelCase(reultList.get(0));
|
||||
//单个Map类型
|
||||
if (resultType==Map.class){
|
||||
return mapHandler;
|
||||
//单个自定义类型
|
||||
} else if (!isTargetClass(resultType)){
|
||||
return BeanConv.toBean(mapHandler, resultType);
|
||||
//单个JDK提供指定类型
|
||||
}else {
|
||||
if (mapHandler.size()!=2){
|
||||
throw new RuntimeException("返回结果非单值");
|
||||
}
|
||||
for (String key : mapHandler.keySet()) {
|
||||
if (!key.equals("time")&&!EmptyUtil.isNullOrEmpty((mapHandler.get(key)))){
|
||||
String target = String.valueOf(mapHandler.get(key)).replace(".0","");
|
||||
return convertStringToObject(resultType,target);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
//当前method声明返回结果为list
|
||||
if (returnTypeTarget== List.class){
|
||||
//驼峰处理
|
||||
List<Map<String, Object>> listHandler = convertKeysToCamelCase(reultList);
|
||||
//list的内部为map结果
|
||||
if (resultType==Map.class){
|
||||
return listHandler;
|
||||
//list的内部为自定义类型
|
||||
}else if (!isTargetClass(resultType)){
|
||||
return BeanConv.toBeanList(listHandler, resultType);
|
||||
//list的内部为JDK提供指定类型
|
||||
}else {
|
||||
List<Object> listResult = Lists.newArrayList();
|
||||
listHandler.forEach(mapHandler->{
|
||||
if (mapHandler.size()!=2){
|
||||
throw new RuntimeException("返回结果非单值");
|
||||
}
|
||||
for (String key : mapHandler.keySet()) {
|
||||
if (!key.equals("time")&&!EmptyUtil.isNullOrEmpty((mapHandler.get(key)))){
|
||||
String target = String.valueOf(mapHandler.get(key)).replace(".0","");
|
||||
listResult.add(convertStringToObject(resultType,target));
|
||||
}
|
||||
}
|
||||
});
|
||||
return listResult;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
// 检查类是否是目标类型
|
||||
public static boolean isTargetClass(Class<?> clazz) {
|
||||
return clazz == Integer.class ||
|
||||
clazz == int.class ||
|
||||
clazz == Long.class ||
|
||||
clazz == long.class ||
|
||||
clazz == Float.class ||
|
||||
clazz == float.class ||
|
||||
clazz == Double.class ||
|
||||
clazz == double.class ||
|
||||
clazz == Short.class ||
|
||||
clazz == short.class ||
|
||||
clazz == Byte.class ||
|
||||
clazz == byte.class ||
|
||||
clazz == Character.class ||
|
||||
clazz == char.class ||
|
||||
clazz == Boolean.class||
|
||||
clazz == boolean.class||
|
||||
clazz== BigDecimal.class ||
|
||||
clazz== String.class;
|
||||
}
|
||||
|
||||
public static Map<String, Object> convertKeysToCamelCase(Map<String, Object> map) {
|
||||
Map<String, Object> camelCaseMap = new HashMap<>();
|
||||
|
||||
for (Map.Entry<String, Object> entry : map.entrySet()) {
|
||||
String originalKey = entry.getKey();
|
||||
Object value = entry.getValue();
|
||||
String camelCaseKey = convertToCamelCase(originalKey);
|
||||
|
||||
camelCaseMap.put(camelCaseKey, value);
|
||||
}
|
||||
|
||||
return camelCaseMap;
|
||||
}
|
||||
|
||||
public static List<Map<String, Object>> convertKeysToCamelCase(List<Map<String, Object>> mapList) {
|
||||
List<Map<String, Object>> listHandler = Lists.newArrayList();
|
||||
mapList.forEach(n->{
|
||||
listHandler.add(convertKeysToCamelCase(n));
|
||||
});
|
||||
return listHandler;
|
||||
}
|
||||
|
||||
public static String convertToCamelCase(String snakeCase) {
|
||||
StringBuilder camelCase = new StringBuilder();
|
||||
boolean nextUpperCase = false;
|
||||
for (int i = 0; i < snakeCase.length(); i++) {
|
||||
char currentChar = snakeCase.charAt(i);
|
||||
if (currentChar == '_') {
|
||||
nextUpperCase = true;
|
||||
} else {
|
||||
if (nextUpperCase) {
|
||||
camelCase.append(Character.toUpperCase(currentChar));
|
||||
nextUpperCase = false;
|
||||
} else {
|
||||
camelCase.append(Character.toLowerCase(currentChar));
|
||||
}
|
||||
}
|
||||
}
|
||||
return camelCase.toString();
|
||||
}
|
||||
|
||||
@SneakyThrows
|
||||
public static <T> T convertStringToObject(Class<?> clazz, String str){
|
||||
if (clazz == String.class) {
|
||||
return (T)str; // 如果目标类型是 String,则直接返回字符串
|
||||
} else if (isTargetClass(clazz)){
|
||||
// 获取目标类型的构造函数,参数为 String 类型的参数
|
||||
Constructor<?> constructor = clazz.getConstructor(String.class);
|
||||
return (T)constructor.newInstance(str); // 使用构造函数创建目标类型的对象
|
||||
}else {
|
||||
return (T)clazz.newInstance();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1 @@
|
||||
org.springframework.boot.autoconfigure.EnableAutoConfiguration=com.itheima.sfbx.framework.influxdb.config.InfluxDBConfig
|
||||
Reference in New Issue
Block a user