Skip to content

Commit e76af6d

Browse files
committed
JdbcDao support including status column in Query Automatic,Stream read and write support five Serialize solution
1 parent ec55967 commit e76af6d

18 files changed

Lines changed: 475 additions & 256 deletions

File tree

common/src/main/java/com/robin/core/fileaccess/iterator/AbstractResIterator.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,7 @@
11
package com.robin.core.fileaccess.iterator;
22

3+
import com.google.gson.Gson;
4+
import com.robin.comm.util.json.GsonUtil;
35
import com.robin.core.base.exception.OperationNotSupportException;
46
import com.robin.core.fileaccess.fs.AbstractFileSystemAccessor;
57
import com.robin.core.fileaccess.meta.DataCollectionMeta;
@@ -22,6 +24,7 @@ public abstract class AbstractResIterator implements IResourceIterator {
2224
protected Map<String, DataSetColumnMeta> columnMap=new HashMap<>();
2325
protected Logger logger= LoggerFactory.getLogger(getClass());
2426
protected String identifier;
27+
protected Gson gson= GsonUtil.getGson();
2528
protected AbstractResIterator(){
2629

2730
}

common/src/main/java/com/robin/core/fileaccess/writer/AbstractQueueWriter.java

Lines changed: 0 additions & 96 deletions
This file was deleted.

common/src/main/java/com/robin/core/fileaccess/writer/AbstractResourceWriter.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,7 @@ public abstract class AbstractResourceWriter implements IResourceWriter{
2525
protected Map<String, Void> columnMap=new HashMap<String, Void>();
2626
protected List<String> columnList=new ArrayList<String>();
2727
protected Logger logger= LoggerFactory.getLogger(getClass());
28-
protected String valueType= ResourceConst.VALUE_TYPE.JSON.getValue();
28+
protected String valueType= ResourceConst.SERIALIZETYPE.JSON.getValue();
2929
protected Schema schema;
3030
protected Map<String,Object> cfgMap;
3131
protected Gson gson= GsonUtil.getGson();

common/src/main/java/com/robin/core/resaccess/iterator/AbstractQueueIterator.java

Lines changed: 0 additions & 59 deletions
This file was deleted.

core/src/main/java/com/robin/core/base/annotation/MappingEntity.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -28,4 +28,6 @@
2828
String jdbcDao() default "jdbcDao";
2929
//Does all field must declare explicit?
3030
boolean explicit() default false;
31+
String statusColumn() default "";
32+
String statusValue() default "";
3133
}

core/src/main/java/com/robin/core/base/dao/JdbcDao.java

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -322,13 +322,19 @@ public <T extends BaseObject> List<T> queryByField(Class<T> type, String fieldNa
322322
StringBuilder buffer = new StringBuilder();
323323
buffer.append(getWholeSelectSql(type)).append(Const.SQL_WHERE);
324324
StringBuilder queryBuffer = new StringBuilder();
325+
AnnotationRetriever.EntityContent<T> tableDef = AnnotationRetriever.getMappingTableByCache(type);
325326
Map<String, FieldContent> map1 = AnnotationRetriever.getMappingFieldsMapCache(type);
326327
List<FieldContent> fields = AnnotationRetriever.getMappingFieldsCache(type);
327328
List<Map<String, Object>> rsList;
328329
if (map1.containsKey(fieldName)) {
330+
List<Object> queryParam=Lists.newArrayList(fieldValues);
329331
generateQuerySqlBySingleFields(map1.get(fieldName), oper, queryBuffer, fieldValues.length);
332+
if(tableDef.isContainStatusColumn()){
333+
appendStatusColumn(queryBuffer,tableDef.getStatusColumn());
334+
queryParam.add(tableDef.getStatusValue());
335+
}
330336
buffer.append(queryBuffer);
331-
rsList = queryBySql(buffer.toString(), fieldValues);
337+
rsList = queryBySql(buffer.toString(), queryParam.toArray());
332338
wrapList(type, retlist, fields, rsList);
333339
} else {
334340
throw new DAOException("query Field not in entity");
@@ -374,19 +380,25 @@ public <T extends BaseObject> List<T> queryByFieldOrderBy(Class<T> type, String
374380
List<T> retlist = new ArrayList<>();
375381
try {
376382
StringBuilder builder = new StringBuilder();
383+
AnnotationRetriever.EntityContent<T> tableDef = AnnotationRetriever.getMappingTableByCache(type);
377384
List<FieldContent> fields = AnnotationRetriever.getMappingFieldsCache(type);
378385
builder.append(getWholeSelectSql(type)).append(Const.SQL_WHERE);
379386
StringBuilder queryBuffer = new StringBuilder();
380387
Map<String, FieldContent> map1 = AnnotationRetriever.getMappingFieldsMapCache(type);
381388

382389
List<Map<String, Object>> rsList;
383390
if (map1.containsKey(fieldName)) {
391+
List<Object> queryParam=Lists.newArrayList(fieldValues);
384392
generateQuerySqlBySingleFields(map1.get(fieldName), oper, queryBuffer, fieldValues.length);
393+
if(tableDef.isContainStatusColumn()){
394+
appendStatusColumn(queryBuffer,tableDef.getStatusColumn());
395+
queryParam.add(tableDef.getStatusValue());
396+
}
385397
builder.append(queryBuffer);
386398
if (!ObjectUtils.isEmpty(orderByStr)) {
387399
builder.append(" order by ").append(orderByStr);
388400
}
389-
rsList = queryBySql(builder.toString(), fieldValues);
401+
rsList = queryBySql(builder.toString(), queryParam.toArray());
390402
wrapList(type, retlist, fields, rsList);
391403
} else {
392404
throw new DAOException("query Field not in entity");
@@ -942,6 +954,9 @@ private void generateQuerySqlBySingleFields(FieldContent columncfg, Const.OPERAT
942954
break;
943955
}
944956
}
957+
private void appendStatusColumn(StringBuilder queryBuilder,String statusColumn){
958+
queryBuilder.append(" AND ").append(statusColumn).append("=?");
959+
}
945960

946961
private void wrapResultToModelWithKey(BaseObject obj, Map<String, Object> map, List<FieldContent> fields, Serializable pkObj) throws Throwable {
947962
for (FieldContent field : fields) {

core/src/main/java/com/robin/core/base/dao/util/AnnotationRetriever.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -212,6 +212,13 @@ private static <T extends BaseObject> EntityContent<T> getMappingTableEntity(Cla
212212
if (!ObjectUtils.isEmpty(jdbcDao)) {
213213
content.setJdbcDao(jdbcDao);
214214
}
215+
if(!ObjectUtils.isEmpty(entity.statusColumn())){
216+
content.setStatusColumn(entity.statusColumn());
217+
content.setContainStatusColumn(true);
218+
}
219+
if(!ObjectUtils.isEmpty(entity.statusValue())){
220+
content.setStatusValue(entity.statusValue());
221+
}
215222
} else {
216223
flag = clazz.isAnnotationPresent(Entity.class);
217224
if (flag) {
@@ -749,6 +756,9 @@ public static class EntityContent<T extends BaseObject> {
749756
private String schema;
750757
private String jdbcDao;
751758
private Class<T> clazz;
759+
private boolean containStatusColumn=false;
760+
private String statusColumn;
761+
private String statusValue=Const.VALID;
752762

753763
public EntityContent(String tableName) {
754764
this.tableName = tableName;

core/src/main/java/com/robin/core/base/dao/util/EntityMappingUtil.java

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -356,7 +356,10 @@ public static <T extends BaseObject> SelectSegment getSelectPkSegment(Class<T> c
356356
segment.getAvailableFields().add(field);
357357
sqlbuffer.append(field.getFieldName()).append(Const.SQL_AS).append(field.getPropertyName()).append(",");
358358
}
359-
359+
}
360+
if(tableDef.isContainStatusColumn()){
361+
wherebuffer.append(tableDef.getStatusColumn()).append("=? and ");
362+
selectObjs.add(tableDef.getStatusValue());
360363
}
361364
sqlbuffer.deleteCharAt(sqlbuffer.length() - 1).append(Const.SQL_FROM);
362365
appendSchemaAndTable(tableDef, sqlbuffer, sqlGen);

core/src/main/java/com/robin/core/base/util/ResourceConst.java

Lines changed: 17 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,9 @@
11
package com.robin.core.base.util;
22

33

4+
import org.springframework.util.Assert;
5+
import org.springframework.util.ObjectUtils;
6+
47
public class ResourceConst {
58
public static final String WORKINGPATHPARAM="output.workingPath";
69
public static final String USETMPFILETAG="output.usingTmpFiles";
@@ -20,6 +23,7 @@ public class ResourceConst {
2023
public static final String STORAGEFILTERSQL="storage.FilterSql";
2124
public static final String PARQUETFILEFORMAT="parquet.file.format";
2225
public static final String USEADMINTG="useAdmin";
26+
public static final String SERIALIZETYPE_COLUMN="output.SerializeType";
2327

2428
public enum IngestType {
2529
TYPE_HDFS(1L,"HDFS"),
@@ -93,22 +97,29 @@ public String toString() {
9397
return String.valueOf(this.value);
9498
}
9599
}
96-
public enum VALUE_TYPE{
100+
public enum SERIALIZETYPE {
97101
AVRO("avro"),
98102
JSON("json"),
99-
XML("xml"),
100103
PROTOBUF("proto"),
101-
ORC("orc"),
102-
CSV("csv"),
103-
PARQUET("parquet");
104+
MESSAGEPACK("msgpack"),
105+
BIJECTION("bijection");
104106
private String value;
105-
VALUE_TYPE(String value){
107+
SERIALIZETYPE(String value){
106108
this.value=value;
107109
}
108110

109111
public String getValue() {
110112
return value;
111113
}
114+
public static SERIALIZETYPE ParseFrom(String input){
115+
Assert.isTrue(!ObjectUtils.isEmpty(input),"");
116+
for(SERIALIZETYPE type :values()){
117+
if(type.getValue().equalsIgnoreCase(input)){
118+
return type;
119+
}
120+
}
121+
throw new IllegalArgumentException("unknown valueType"+input);
122+
}
112123
}
113124
public enum RESTYPE{
114125
DIR("1"),

core/src/main/java/com/robin/core/fileaccess/meta/DataCollectionMeta.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -126,6 +126,9 @@ public boolean isQueueType(){
126126
}
127127
public static class Builder{
128128
private final DataCollectionMeta meta=new DataCollectionMeta();
129+
public static Builder newBuilder(){
130+
return new Builder();
131+
}
129132
public Builder(){
130133

131134
}

0 commit comments

Comments
 (0)