7c237766460e48c0d6c336af4f26de253f561112
[dcaegen2/services.git] /
1 /*
2 * ============LICENSE_START=======================================================
3 * ONAP : DATALAKE
4 * ================================================================================
5 * Copyright 2019 China Mobile
6 *=================================================================================
7 * Licensed under the Apache License, Version 2.0 (the "License");
8 * you may not use this file except in compliance with the License.
9 * You may obtain a copy of the License at
10 *
11 *     http://www.apache.org/licenses/LICENSE-2.0
12 *
13 * Unless required by applicable law or agreed to in writing, software
14 * distributed under the License is distributed on an "AS IS" BASIS,
15 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
16 * See the License for the specific language governing permissions and
17 * limitations under the License.
18 * ============LICENSE_END=========================================================
19 */
20
21 package org.onap.datalake.feeder.service;
22
23 import java.util.ArrayList;
24 import java.util.List;
25
26 import javax.annotation.PostConstruct;
27 import javax.annotation.PreDestroy;
28
29 import org.json.JSONObject;
30 import org.onap.datalake.feeder.config.ApplicationConfiguration;
31 import org.onap.datalake.feeder.domain.Db;
32 import org.onap.datalake.feeder.dto.TopicConfig;
33 import org.slf4j.Logger;
34 import org.slf4j.LoggerFactory;
35
36 import org.springframework.beans.factory.annotation.Autowired;
37 import org.springframework.stereotype.Service;
38
39 import com.couchbase.client.java.Bucket;
40 import com.couchbase.client.java.Cluster;
41 import com.couchbase.client.java.CouchbaseCluster;
42 import com.couchbase.client.java.document.JsonDocument;
43 import com.couchbase.client.java.document.JsonLongDocument;
44 import com.couchbase.client.java.document.json.JsonObject;
45 import com.couchbase.client.java.env.CouchbaseEnvironment;
46 import com.couchbase.client.java.env.DefaultCouchbaseEnvironment;
47
48 import rx.Observable;
49 import rx.functions.Func1;
50
51 /**
52  * Service to use Couchbase
53  * 
54  * @author Guobiao Mo
55  *
56  */
57 @Service
58 public class CouchbaseService {
59
60         private final Logger log = LoggerFactory.getLogger(this.getClass());
61
62         @Autowired
63         ApplicationConfiguration config;
64
65         @Autowired
66         private DbService dbService;
67
68         Bucket bucket;
69         private boolean isReady = false;
70
71         @PostConstruct
72         private void init() {
73                 // Initialize Couchbase Connection
74                 try {
75                         Db couchbase = dbService.getCouchbase();
76
77                         //this tunes the SDK (to customize connection timeout)
78                         CouchbaseEnvironment env = DefaultCouchbaseEnvironment.builder().connectTimeout(60000) // 60s, default is 5s
79                                         .build();
80                         Cluster cluster = CouchbaseCluster.create(env, couchbase.getHost());
81                         cluster.authenticate(couchbase.getLogin(), couchbase.getPass());
82                         bucket = cluster.openBucket(couchbase.getDatabase());
83                         // Create a N1QL Primary Index (but ignore if it exists)
84                         bucket.bucketManager().createN1qlPrimaryIndex(true, false);
85
86                         log.info("Connected to Couchbase {}", couchbase.getHost());
87                         isReady = true;
88                 } catch (Exception ex) {
89                         log.error("error connection to Couchbase.", ex);
90                         isReady = false;
91                 }
92         }
93
94         @PreDestroy
95         public void cleanUp() {
96                 bucket.close();
97         }
98
99         public void saveJsons(TopicConfig topic, List<JSONObject> jsons) {
100                 List<JsonDocument> documents = new ArrayList<>(jsons.size());
101                 for (JSONObject json : jsons) {
102                         //convert to Couchbase JsonObject from org.json JSONObject
103                         JsonObject jsonObject = JsonObject.fromJson(json.toString());
104
105                         long timestamp = jsonObject.getLong(config.getTimestampLabel());//this is Kafka time stamp, which is added in StoreService.messageToJson()
106
107                         //setup TTL
108                         int expiry = (int) (timestamp / 1000L) + topic.getTtl() * 3600 * 24; //in second
109
110                         String id = getId(topic, json);
111                         JsonDocument doc = JsonDocument.create(id, expiry, jsonObject);
112                         documents.add(doc);
113                 }
114                 try {
115                         saveDocuments(documents);
116                 } catch (Exception e) {
117                         log.error("error saving to Couchbase.", e);
118                 }
119                 log.debug("saved text to topic = {}, this batch count = {} ", topic, documents.size());
120         }
121
122         public String getId(TopicConfig topic, JSONObject json) {
123                 //if this topic requires extract id from JSON
124                 String id = topic.getMessageId(json);
125                 if (id != null) {
126                         return id;
127                 }
128
129                 String topicStr = topic.getName();
130                 //String id = topicStr+":"+timestamp+":"+UUID.randomUUID();
131
132                 //https://forums.couchbase.com/t/how-to-set-an-auto-increment-id/4892/2
133                 //atomically get the next sequence number:
134                 // increment by 1, initialize at 0 if counter doc not found
135                 //TODO how slow is this compared with above UUID approach?
136                 JsonLongDocument nextIdNumber = bucket.counter(topicStr, 1, 0); //like 12345 
137                 id = topicStr + ":" + nextIdNumber.content();
138
139                 return id;
140         }
141
142         //https://docs.couchbase.com/java-sdk/2.7/document-operations.html
143         private void saveDocuments(List<JsonDocument> documents) {
144                 Observable.from(documents).flatMap(new Func1<JsonDocument, Observable<JsonDocument>>() {
145                         @Override
146                         public Observable<JsonDocument> call(final JsonDocument docToInsert) {
147                                 return bucket.async().insert(docToInsert);
148                         }
149                 }).last().toBlocking().single();
150         }
151
152 }