Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 9 additions & 9 deletions jobs-core/src/main/scala/org/sunbird/job/util/Neo4JUtil.scala
Original file line number Diff line number Diff line change
@@ -1,17 +1,17 @@
package org.sunbird.job.util

import java.util

import org.neo4j.driver.v1.{Config, GraphDatabase}
import org.neo4j.driver.v1.{Config, Driver, GraphDatabase, StatementResult}
import org.slf4j.LoggerFactory

import java.util
import scala.collection.JavaConverters._

class Neo4JUtil(routePath: String, graphId: String) {
class Neo4JUtil(routePath: String, graphId: String) extends Serializable {

private[this] val logger = LoggerFactory.getLogger(classOf[Neo4JUtil])

val maxIdleSession = 20
val driver = GraphDatabase.driver(routePath, getConfig)
val driver: Driver = GraphDatabase.driver(routePath, getConfig)

def getConfig: Config = {
val config = Config.build
Expand All @@ -33,7 +33,7 @@ class Neo4JUtil(routePath: String, graphId: String) {

def getNodeProperties(identifier: String): java.util.Map[String, AnyRef] = {
val session = driver.session()
val query = s"""MATCH (n:${graphId}{IL_UNIQUE_ID:"${identifier}"}) return n;"""
val query = s"""MATCH (n:$graphId{IL_UNIQUE_ID:"$identifier"}) return n;"""
val statementResult = session.run(query)
if (statementResult.hasNext)
statementResult.single().get("n").asMap()
Expand All @@ -42,7 +42,7 @@ class Neo4JUtil(routePath: String, graphId: String) {

def getNodePropertiesWithObjectType(objectType: String): util.List[util.Map[String, AnyRef]] = {
val session = driver.session()
val query = s"""MATCH (n:${graphId}) where n.IL_FUNC_OBJECT_TYPE = "${objectType}" AND n.IL_SYS_NODE_TYPE="DATA_NODE" return n;"""
val query = s"""MATCH (n:$graphId) where n.IL_FUNC_OBJECT_TYPE = "$objectType" AND n.IL_SYS_NODE_TYPE="DATA_NODE" return n;"""
val statementResult = session.run(query)
if (statementResult.hasNext)
statementResult.list().asScala.toList.map(record => record.get("n").asMap()).asJava
Expand Down Expand Up @@ -71,7 +71,7 @@ class Neo4JUtil(routePath: String, graphId: String) {
else throw new Exception(s"Unable to update the node with identifier: $identifier")
}

def executeQuery(query: String) = {
def executeQuery(query: String): StatementResult = {
val session = driver.session()
session.run(query)
}
Expand All @@ -83,7 +83,7 @@ class Neo4JUtil(routePath: String, graphId: String) {

def updateNode(identifier: String, metadata: Map[String, AnyRef]): Unit = {
val query = s"""MATCH (n:$graphId {IL_UNIQUE_ID:"$identifier"}) SET n = """ + "$properties return n;"
logger.info(s"Query for updating metadata for identifier : ${identifier} is : ${query}")
logger.info(s"Query for updating metadata for identifier : $identifier is : $query")
val session = driver.session()
val properties: java.util.Map[String, AnyRef] = Map[String, AnyRef]("properties" -> metadata.asJava).asJava
val result = session.run(query, properties)
Expand Down
6 changes: 6 additions & 0 deletions publish-pipeline/content-publish/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,12 @@
<version>3.11.2.0</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>it.ozimov</groupId>
<artifactId>embedded-redis</artifactId>
<version>0.7.3</version>
<scope>test</scope>
</dependency>
</dependencies>

<build>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ import java.util.UUID
import scala.concurrent.ExecutionContext

class ContentPublishFunction(config: ContentPublishConfig, httpUtil: HttpUtil,
@transient var neo4JUtil: Neo4JUtil = null,
neo4JUtil: Neo4JUtil,
@transient var cassandraUtil: CassandraUtil = null,
@transient var cloudStorageUtil: CloudStorageUtil = null,
@transient var definitionCache: DefinitionCache = null,
Expand All @@ -43,7 +43,6 @@ class ContentPublishFunction(config: ContentPublishConfig, httpUtil: HttpUtil,
override def open(parameters: Configuration): Unit = {
super.open(parameters)
cassandraUtil = new CassandraUtil(config.cassandraHost, config.cassandraPort)
neo4JUtil = new Neo4JUtil(config.graphRoutePath, config.graphName)
cloudStorageUtil = new CloudStorageUtil(config)
ec = ExecutionContexts.global
definitionCache = new DefinitionCache()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,12 +8,12 @@ import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
import org.sunbird.job.connector.FlinkKafkaConnector
import org.sunbird.job.content.function.{ContentPublishFunction, PublishEventRouter}
import org.sunbird.job.content.publish.domain.Event
import org.sunbird.job.util.{FlinkUtil, HttpUtil}
import org.sunbird.job.util.{FlinkUtil, HttpUtil, Neo4JUtil}

import java.io.File
import java.util

class ContentPublishStreamTask(config: ContentPublishConfig, kafkaConnector: FlinkKafkaConnector, httpUtil: HttpUtil) {
class ContentPublishStreamTask(config: ContentPublishConfig, kafkaConnector: FlinkKafkaConnector, httpUtil: HttpUtil, neo4JUtil: Neo4JUtil) {

def process(): Unit = {
implicit val env: StreamExecutionEnvironment = FlinkUtil.getExecutionContext(config)
Expand All @@ -29,7 +29,7 @@ class ContentPublishStreamTask(config: ContentPublishConfig, kafkaConnector: Fli
.name("publish-event-router").uid("publish-event-router")
.setParallelism(config.eventRouterParallelism)

val contentPublish = processStreamTask.getSideOutput(config.contentPublishOutTag).process(new ContentPublishFunction(config, httpUtil))
val contentPublish = processStreamTask.getSideOutput(config.contentPublishOutTag).process(new ContentPublishFunction(config, httpUtil, neo4JUtil))
.name("content-publish-process").uid("content-publish-process").setParallelism(1)

contentPublish.getSideOutput(config.generateVideoStreamingOutTag).addSink(kafkaConnector.kafkaStringSink(config.postPublishTopic))
Expand All @@ -52,7 +52,8 @@ object ContentPublishStreamTask {
val publishConfig = new ContentPublishConfig(config)
val kafkaUtil = new FlinkKafkaConnector(publishConfig)
val httpUtil = new HttpUtil
val task = new ContentPublishStreamTask(publishConfig, kafkaUtil, httpUtil)
val neo4JUtil = new Neo4JUtil(publishConfig.graphRoutePath, publishConfig.graphName)
val task = new ContentPublishStreamTask(publishConfig, kafkaUtil, httpUtil, neo4JUtil)
task.process()
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import com.typesafe.config.{Config, ConfigFactory}
import org.apache.flink.api.common.typeinfo.TypeInformation
import org.apache.flink.api.java.typeutils.TypeExtractor
import org.apache.flink.runtime.testutils.MiniClusterResourceConfiguration
import org.apache.flink.streaming.api.functions.sink.SinkFunction
import org.apache.flink.streaming.api.functions.source.SourceFunction
import org.apache.flink.streaming.api.functions.source.SourceFunction.SourceContext
import org.apache.flink.test.util.MiniClusterWithClientResource
Expand All @@ -14,13 +15,16 @@ import org.cassandraunit.utils.EmbeddedCassandraServerHelper
import org.mockito.ArgumentMatchers.anyString
import org.mockito.Mockito
import org.mockito.Mockito.when
import org.sunbird.job.cache.{DataCache, RedisConnect}
import org.sunbird.job.connector.FlinkKafkaConnector
import org.sunbird.job.content.publish.domain.Event
import org.sunbird.job.content.task.{ContentPublishConfig, ContentPublishStreamTask}
import org.sunbird.job.fixture.EventFixture
import org.sunbird.job.publish.config.PublishConfig
import org.sunbird.job.util.{CassandraUtil, CloudStorageUtil, HttpUtil, Neo4JUtil}
import org.sunbird.spec.{BaseMetricsReporter, BaseTestSpec}
import redis.clients.jedis.Jedis
import redis.embedded.RedisServer

import java.util

Expand All @@ -38,19 +42,28 @@ class ContentPublishStreamTaskSpec extends BaseTestSpec {
val config: Config = ConfigFactory.load("test.conf").withFallback(ConfigFactory.systemEnvironment())
implicit val jobConfig: ContentPublishConfig = new ContentPublishConfig(config)

val mockHttpUtil = mock[HttpUtil](Mockito.withSettings().serializable())
implicit val mockNeo4JUtil: Neo4JUtil = mock[Neo4JUtil](Mockito.withSettings().serializable())
// val mockHttpUtil: HttpUtil = mock[HttpUtil](Mockito.withSettings().serializable())
val mockHttpUtil: HttpUtil = new HttpUtil
val mockNeo4JUtil: Neo4JUtil = mock[Neo4JUtil](Mockito.withSettings().serializable())
var cassandraUtil: CassandraUtil = _
val publishConfig: PublishConfig = new PublishConfig(config, "")
val cloudStorageUtil: CloudStorageUtil = new CloudStorageUtil(publishConfig)
var jedis: Jedis = _
val mockDataCache: DataCache = mock[DataCache](Mockito.withSettings().serializable())
var redisServer: RedisServer = _

override protected def beforeAll(): Unit = {
super.beforeAll()
redisServer = new RedisServer(6340)
redisServer.start()
val redisConnect = new RedisConnect(jobConfig)
jedis = redisConnect.getConnection(jobConfig.nodeStore)
EmbeddedCassandraServerHelper.startEmbeddedCassandra(80000L)
cassandraUtil = new CassandraUtil(jobConfig.cassandraHost, jobConfig.cassandraPort)
val session = cassandraUtil.session
val dataLoader = new CQLDataLoader(session)
dataLoader.load(new FileCQLDataSet(getClass.getResource("/test.cql").getPath, true, true))
jedis.flushDB()
flinkCluster.before()
}

Expand All @@ -59,32 +72,34 @@ class ContentPublishStreamTaskSpec extends BaseTestSpec {
try {
EmbeddedCassandraServerHelper.cleanEmbeddedCassandra()
} catch {
case ex: Exception => {
}
case ex: Exception =>
}
redisServer.stop()
flinkCluster.after()
}

def initialize(): Unit = {
when(mockKafkaUtil.kafkaJobRequestSource[Event](jobConfig.kafkaInputTopic)).thenReturn(new ContentPublishEventSource)
when(mockKafkaUtil.kafkaStringSink(jobConfig.postPublishTopic)).thenReturn(new PostPublishEventSink)
when(mockKafkaUtil.kafkaStringSink(jobConfig.kafkaErrorTopic)).thenReturn(new ContentPublishFailedEventSink)
}

ignore should " publish the content " in {
"ContentPublishStreamTask process event " should " publish the content " in {
when(mockNeo4JUtil.getNodeProperties(anyString())).thenReturn(new util.HashMap[String, AnyRef])
initialize
new ContentPublishStreamTask(jobConfig, mockKafkaUtil, mockHttpUtil).process()
initialize()
new ContentPublishStreamTask(jobConfig, mockKafkaUtil, mockHttpUtil, mockNeo4JUtil).process()
BaseMetricsReporter.gaugeMetrics(s"${jobConfig.jobName}.${jobConfig.totalEventsCount}").getValue() should be(1)
BaseMetricsReporter.gaugeMetrics(s"${jobConfig.jobName}.${jobConfig.contentPublishEventCount}").getValue() should be(1)
}
}

private class ContentPublishEventSource extends SourceFunction[Event] {

override def run(ctx: SourceContext[Event]) {
override def run(ctx: SourceContext[Event]): Unit = {
ctx.collect(jsonToEvent(EventFixture.PDF_EVENT1))
}

override def cancel() = {}
override def cancel(): Unit = {}

def jsonToEvent(json: String): Event = {
val gson = new Gson()
Expand All @@ -94,3 +109,29 @@ private class ContentPublishEventSource extends SourceFunction[Event] {
new Event(data, 0, 10)
}
}

class ContentPublishFailedEventSink extends SinkFunction[String] {

override def invoke(value: String): Unit = {
synchronized {
ContentPublishFailedEventSink.values.add(value)
}
}
}

object ContentPublishFailedEventSink {
val values: util.List[String] = new util.ArrayList()
}

class PostPublishEventSink extends SinkFunction[String] {

override def invoke(value: String): Unit = {
synchronized {
PostPublishEventSink.values.add(value)
}
}
}

object PostPublishEventSink {
val values: util.List[String] = new util.ArrayList()
}