Skip to content

Commit 12dd689

Browse files
committed
add HTML stripping and HTML entity unescaping custom UDF's and a model for adding UDF's internally in general.
fixed state handling to be index+statement-id specific so that multiple indexes could be managed separately from same statements release 0.6.0-ALPHA
1 parent 377fe9b commit 12dd689

8 files changed

Lines changed: 133 additions & 49 deletions

File tree

README.md

Lines changed: 11 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -59,9 +59,7 @@ Lastly you specify `importSteps` which are queries that use any of the `sources`
5959
"username": "myusername",
6060
"password": "mypass"
6161
},
62-
"driverJars": [
63-
"/Users/myUserName/Downloads/mysql-connector-java-5.1.41/mysql-connector-java-5.1.41-bin.jar"
64-
],
62+
"driverJars": [], # MySQL and Postgres JARS are included automatically, this property can be omitted completely
6563
"tables": [
6664
{
6765
"sparkTable": "Users",
@@ -367,13 +365,21 @@ for full details on the supported SQL.
367365
A list of [SQL Functions](https://spark.apache.org/docs/2.1.0/api/java/org/apache/spark/sql/functions.html)
368366
is available in raw API docs. (_TODO: find better reference_)
369367

368+
The Data Import Handler also defines some UDF functions for use within SQL:
369+
370+
|Function|Description|
371+
|-------|-----------|
372+
|stripHtml(string)|Removes all HTML tags and returns only the text (including unescaping of HTML Entities)|
373+
|unescapeHtmlEntites(string)|Unescapes HTML entities found in the text|
374+
|fluffly(string)|A silly function that prepends the word "fluffly" to the text, used as a text function to mark values as being changed by processing|
375+
370376
### State Management and History:
371377

372378
State for the `lastRun` value is per-statement and stored in the target Elasticsearch cluster for that statement. An index
373-
will be created called `.kohesive-dih-state` which stores the last run state, a lock for current running statements, and
379+
will be created called `.kohesive-dih-state-v2` which stores the last run state, a lock for current running statements, and
374380
a log of all previous runs (success and failures, along with row count processed by the statement query).
375381

376-
You should inspect this log (index `.kohesive-dih-state` type `log`)if you wish to monitor the results of runs.
382+
You should inspect this log (index `.kohesive-dih-state-v2` type `log`)if you wish to monitor the results of runs.
377383

378384
### Parallelism
379385

build.gradle

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,8 @@ dependencies {
3535

3636
compile group: 'org.apache.spark', name: "spark-sql_${version_scala}", version: '2.1.0'
3737

38+
compile group: 'org.jsoup', name: 'jsoup', version: version_jsoup
39+
3840
resolutionRules group: 'com.netflix.nebula', name: 'gradle-resolution-rules', version: version_nebula_resolution
3941
resolutionRules files("${rootDir}/gradle-local-dependency-rules.json")
4042

gradle.properties

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,8 @@ version_elasticsearch=5.2.2
1414
version_okhttp=3.6.0
1515
version_scala=2.11
1616

17+
version_jsoup=1.10.2
18+
1719
version_slf4j=1.+
1820
version_logback=1.+
1921
version_junit=4.+
Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,22 @@
1+
package uy.kohesive.elasticsearch.dataimport.udf;
2+
3+
import org.apache.spark.sql.SparkSession;
4+
import org.apache.spark.sql.api.java.UDF1;
5+
import org.apache.spark.sql.api.java.UDF2;
6+
import org.apache.spark.sql.types.DataType;
7+
import org.apache.spark.sql.types.DataTypes;
8+
import org.jsoup.Jsoup;
9+
import org.jsoup.nodes.Element;
10+
import org.jsoup.safety.Whitelist;
11+
12+
public class Udfs {
13+
public static <T extends String, RT extends String> void registerStringToStringUdf(final SparkSession spark, final String name, final UDF1<T, RT> f) {
14+
spark.udf().register(name, f, DataTypes.StringType);
15+
}
16+
17+
public static <T1 extends String, T2 extends String, RT extends String> void registerStringStringToStringUdf(final SparkSession spark, final String name, final UDF2<T1, T2, RT> f) {
18+
spark.udf().register(name, f, DataTypes.StringType);
19+
}
20+
21+
22+
}

src/main/kotlin/uy/kohesive/elasticsearch/dataimport/App.kt

Lines changed: 17 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import org.apache.spark.sql.AnalysisException
1010
import org.apache.spark.sql.SparkSession
1111
import org.apache.spark.sql.catalyst.parser.ParseException
1212
import org.elasticsearch.spark.sql.api.java.JavaEsSparkSQL
13+
import uy.kohesive.elasticsearch.dataimport.udf.Udfs
1314
import java.io.File
1415
import java.io.InputStream
1516
import java.io.InputStreamReader
@@ -71,9 +72,9 @@ class App {
7172

7273
val lastRuns: Map<String, Instant> = cfg.importSteps.map { importStep ->
7374
importStep.statements.map { statement ->
74-
val lastState = stateMap.get(statement.id)!!.readStateForStatement(uniqueId, statement.id)?.truncatedTo(ChronoUnit.SECONDS) ?: NOSTATE
75+
val lastState = stateMap.get(statement.id)!!.readStateForStatement(uniqueId, statement)?.truncatedTo(ChronoUnit.SECONDS) ?: NOSTATE
7576
println(" Statement ${statement.id} - ${statement.description}")
76-
println(" LAST RUN: ${if (lastState === NOSTATE) "never" else lastState.toIsoString()}")
77+
println(" LAST RUN: ${if (lastState == NOSTATE) "never" else lastState.toIsoString()}")
7778
statement.id to lastState
7879
}
7980
}.flatten().toMap()
@@ -92,9 +93,7 @@ class App {
9293
.master(sparkMaster).getOrCreate().use { spark ->
9394

9495
// add extra UDF functions
95-
// TODO: udf's
96-
97-
// spark.udf().register("fluffy")
96+
DataImportHandlerUdfs.registerSparkUdfs(spark)
9897

9998
// setup FILE inputs
10099
println()
@@ -206,8 +205,13 @@ class App {
206205
val sqlMinDate = Timestamp.from(lastRun).toString()
207206
val sqlMaxDate = Timestamp.from(thisRunDate).toString()
208207

209-
println("\n Execute statement: (range '$sqlMinDate' to '$sqlMaxDate')\n${statement.description.replaceIndent(" ")}")
210-
if (!stateMgr.lockStatement(uniqueId, statement.id)) {
208+
val dateMsg = if (lastRun == NOSTATE) {
209+
"range NEVER to '$sqlMaxDate'"
210+
} else {
211+
"range '$sqlMinDate' to '$sqlMaxDate'"
212+
}
213+
println("\n Execute statement: ($dateMsg)\n${statement.description.replaceIndent(" ")}")
214+
if (!stateMgr.lockStatement(uniqueId, statement)) {
211215
System.err.println(" Cannot aquire lock for statement ${statement.id}")
212216
} else {
213217
try {
@@ -231,14 +235,15 @@ class App {
231235
val rowCount = sqlResults.count()
232236
println(" Rows processed: $rowCount")
233237

234-
stateMgr.writeStateForStatement(uniqueId, statement.id, thisRunDate, "success", rowCount, null)
235-
stateMgr.logStatement(uniqueId, statement.id, thisRunDate, "sucess", rowCount, null)
238+
stateMgr.writeStateForStatement(uniqueId, statement, thisRunDate, "success", rowCount, null)
239+
stateMgr.logStatement(uniqueId, statement, thisRunDate, "sucess", rowCount, null)
236240
} catch (ex: Exception) {
237241
val msg = ex.message ?: "unknown failure"
238-
stateMgr.writeStateForStatement(uniqueId, statement.id, thisRunDate, "error", 0, msg)
239-
stateMgr.logStatement(uniqueId, statement.id, thisRunDate, "error", 0, msg)
242+
stateMgr.writeStateForStatement(uniqueId, statement, lastRun, "error", 0, msg)
243+
stateMgr.logStatement(uniqueId, statement, thisRunDate, "error", 0, msg)
244+
System.err.println("\nProcess FAILED: \n$msg\n")
240245
} finally {
241-
stateMgr.unlockStatemnt(uniqueId, statement.id)
246+
stateMgr.unlockStatemnt(uniqueId, statement)
242247
}
243248
}
244249
}

src/main/kotlin/uy/kohesive/elasticsearch/dataimport/State.kt

Lines changed: 39 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -14,14 +14,16 @@ import java.time.Instant
1414

1515
interface StateManager {
1616
fun init()
17-
fun lockStatement(runId: String, statementId: String): Boolean
18-
fun pingLockStatement(runId: String, statementId: String): Boolean
19-
fun unlockStatemnt(runId: String, statementId: String)
20-
fun writeStateForStatement(runId: String, statementId: String, lastRunStart: Instant, status: String, lastRowCount: Long, errMsg: String? = null)
21-
fun readStateForStatement(runId: String, statementId: String): Instant?
22-
fun logStatement(runId: String, statementId: String, lastRunStart: Instant, status: String, rowCount: Long, errMsg: String? = null)
17+
fun lockStatement(runId: String, statement: EsImportStatement): Boolean
18+
fun pingLockStatement(runId: String, statement: EsImportStatement): Boolean
19+
fun unlockStatemnt(runId: String, statement: EsImportStatement)
20+
fun writeStateForStatement(runId: String, statement: EsImportStatement, lastRunStart: Instant, status: String, lastRowCount: Long, errMsg: String? = null)
21+
fun readStateForStatement(runId: String, statement: EsImportStatement): Instant?
22+
fun logStatement(runId: String, statement: EsImportStatement, lastRunStart: Instant, status: String, rowCount: Long, errMsg: String? = null)
2323
}
2424

25+
private fun EsImportStatement.stateKey(): String = this.indexName + "-" + this.id
26+
2527
// TODO: better state management
2628
// This is NOT using the ES client because we do not want conflicts with Spark dependencies
2729
class ElasticSearchStateManager(val nodes: List<String>, val auth: AuthInfo?) : StateManager {
@@ -32,7 +34,7 @@ class ElasticSearchStateManager(val nodes: List<String>, val auth: AuthInfo?) :
3234
configure(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS, true)
3335
configure(DeserializationFeature.READ_DATE_TIMESTAMPS_AS_NANOSECONDS, false)
3436
}
35-
val STATE_INDEX = ".kohesive-dih-state"
37+
val STATE_INDEX = ".kohesive-dih-state-v2"
3638

3739
private fun OkHttpClient.get(url: String): Pair<Int, String> {
3840
val request = Request.Builder().url(url).build()
@@ -98,6 +100,7 @@ class ElasticSearchStateManager(val nodes: List<String>, val auth: AuthInfo?) :
98100
"mappings": {
99101
"state": {
100102
"properties": {
103+
"targetIndex": { "type": "keyword" },
101104
"statementId": { "type": "keyword" },
102105
"lastRunDate": { "type": "date" },
103106
"status": { "type": "keyword" },
@@ -108,6 +111,7 @@ class ElasticSearchStateManager(val nodes: List<String>, val auth: AuthInfo?) :
108111
},
109112
"log": {
110113
"properties": {
114+
"targetIndex": { "type": "keyword" },
111115
"statementId": { "type": "keyword" },
112116
"runId": { "type": "keyword" },
113117
"runDate": { "type": "date" },
@@ -118,6 +122,7 @@ class ElasticSearchStateManager(val nodes: List<String>, val auth: AuthInfo?) :
118122
},
119123
"lock": {
120124
"properties": {
125+
"targetIndex": { "type": "keyword" },
121126
"statementId": { "type": "keyword" },
122127
"runId": { "type": "keyword" },
123128
"lockDate": { "type": "date" }
@@ -137,7 +142,7 @@ class ElasticSearchStateManager(val nodes: List<String>, val auth: AuthInfo?) :
137142
ttlKillOldLocks()
138143
}
139144

140-
data class Lock(val runId: String, val statementId: String, val lockDate: Instant)
145+
data class Lock(val runId: String, val targetIndex: String, val statementId: String, val lockDate: Instant)
141146

142147
fun makeUrl(index: String, type: String) = "$url/$index/$type"
143148
fun makeUrl(type: String) = makeUrl(STATE_INDEX, type)
@@ -151,19 +156,19 @@ class ElasticSearchStateManager(val nodes: List<String>, val auth: AuthInfo?) :
151156
if (!delCode.isSuccess()) throw DataImportException("State manager failed, TTL delete query for locks failed\n$response")
152157
}
153158

154-
override fun lockStatement(runId: String, statementId: String): Boolean {
155-
val lockUrl = "${makeUrl("lock")}/${statementId}"
156-
val (code, response) = http.post("$lockUrl?op_type=create", JSON.writeValueAsString(Lock(runId, statementId, Instant.now())))
159+
override fun lockStatement(runId: String, statement: EsImportStatement): Boolean {
160+
val lockUrl = "${makeUrl("lock")}/${statement.stateKey()}"
161+
val (code, response) = http.post("$lockUrl?op_type=create", JSON.writeValueAsString(Lock(runId, statement.indexName, statement.id, Instant.now())))
157162

158163
if (!code.isSuccess()) {
159164
ttlKillOldLocks()
160-
return pingLockStatement(runId, statementId)
165+
return pingLockStatement(runId, statement)
161166
}
162167
return true
163168
}
164169

165-
override fun pingLockStatement(runId: String, statementId: String): Boolean {
166-
val lockUrl = "${makeUrl("lock")}/${statementId}"
170+
override fun pingLockStatement(runId: String, statement: EsImportStatement): Boolean {
171+
val lockUrl = "${makeUrl("lock")}/${statement.stateKey()}"
167172
val (getCode, response) = http.get(lockUrl)
168173
if (getCode.isSuccess()) {
169174
val lock = mapFromSource<Lock>(response)
@@ -179,45 +184,47 @@ class ElasticSearchStateManager(val nodes: List<String>, val auth: AuthInfo?) :
179184
"bool": {
180185
"must": [
181186
{ "term": { "runId": "$runId" } },
182-
{ "term": { "statementId": "$statementId" } }
187+
{ "term": { "targetIndex": "${statement.indexName}" } },
188+
{ "term": { "statementId": "${statement.id}" } }
183189
]
184190
}
185191
}
186192
}
187193
""")
188194
if (!updCode.isSuccess()) {
189-
throw DataImportException("State manager failed, cannot acquire lock for $statementId - had conflict on pinging of lock\n$updResponse")
195+
throw DataImportException("State manager failed, cannot acquire lock for ${statement.stateKey()} - had conflict on pinging of lock\n$updResponse")
190196
}
191197
} else {
192-
throw DataImportException("State manager failed, cannot acquire lock for $statementId -- it is held by ${lock.runId} since ${lock.lockDate.toIsoString()}\n$response")
198+
throw DataImportException("State manager failed, cannot acquire lock for ${statement.stateKey()} -- it is held by ${lock.runId} since ${lock.lockDate.toIsoString()}\n$response")
193199
}
194200
}
195201
return true
196202
}
197203

198-
override fun unlockStatemnt(runId: String, statementId: String) {
199-
if (pingLockStatement(runId, statementId)) {
200-
val lockUrl = "${makeUrl("lock")}/${statementId}?refresh"
204+
override fun unlockStatemnt(runId: String, statement: EsImportStatement) {
205+
if (pingLockStatement(runId, statement)) {
206+
val lockUrl = "${makeUrl("lock")}/${statement.stateKey()}?refresh"
201207
val (code, response) = http.delete(lockUrl)
202208
if (!code.isSuccess()) {
203-
throw DataImportException("State manager failed, cannot delete lock for $statementId\n$response")
209+
throw DataImportException("State manager failed, cannot delete lock for ${statement.stateKey()}\n$response")
204210
}
205211
}
206212
}
207213

208-
data class State(val statementId: String, val lastRunDate: Instant, val status: String, val lastRunId: String, val lastErrorMesasge: String?, val lastRowCount: Long)
214+
data class State(val targetIndex: String, val statementId: String, val lastRunDate: Instant, val status: String, val lastRunId: String, val lastErrorMesasge: String?, val lastRowCount: Long)
209215

210-
data class StateLog(val statementId: String, val runId: String, val runDate: Instant, val status: String, val errorMsg: String?, val rowCount: Long)
216+
data class StateLog(val targetIndex: String, val statementId: String, val runId: String, val runDate: Instant, val status: String, val errorMsg: String?, val rowCount: Long)
211217

212-
override fun writeStateForStatement(runId: String, statementId: String, lastRunStart: Instant, status: String, lastRowCount: Long, errMsg: String?) {
213-
val (code, response) = http.post("${makeUrl("state")}/${statementId}?refresh", JSON.writeValueAsString(State(statementId, lastRunStart, status, runId, errMsg, lastRowCount)))
218+
override fun writeStateForStatement(runId: String, statement: EsImportStatement, lastRunStart: Instant, status: String, lastRowCount: Long, errMsg: String?) {
219+
val (code, response) = http.post("${makeUrl("state")}/${statement.stateKey()}?refresh",
220+
JSON.writeValueAsString(State(statement.indexName, statement.id, lastRunStart, status, runId, errMsg, lastRowCount)))
214221
if (!code.isSuccess()) {
215-
throw DataImportException("State manager failed, cannot update state for $statementId\n$response")
222+
throw DataImportException("State manager failed, cannot update state for ${statement.stateKey()}\n$response")
216223
}
217224
}
218225

219-
override fun readStateForStatement(runId: String, statementId: String): Instant? {
220-
val (code, response) = http.get("${makeUrl("state")}/${statementId}")
226+
override fun readStateForStatement(runId: String, statement: EsImportStatement): Instant? {
227+
val (code, response) = http.get("${makeUrl("state")}/${statement.stateKey()}")
221228
if (code.isSuccess()) {
222229
val state = mapFromSource<State>(response)
223230
return state.lastRunDate
@@ -226,10 +233,11 @@ class ElasticSearchStateManager(val nodes: List<String>, val auth: AuthInfo?) :
226233
}
227234
}
228235

229-
override fun logStatement(runId: String, statementId: String, lastRunStart: Instant, status: String, rowCount: Long, errMsg: String?) {
230-
val (code, response) = http.post("${makeUrl("log")}/${statementId}_run_${runId}?refresh", JSON.writeValueAsString(StateLog(statementId, runId, lastRunStart, status, errMsg, rowCount)))
236+
override fun logStatement(runId: String, statement: EsImportStatement, lastRunStart: Instant, status: String, rowCount: Long, errMsg: String?) {
237+
val (code, response) = http.post("${makeUrl("log")}/${statement.stateKey()}_run_${runId}?refresh",
238+
JSON.writeValueAsString(StateLog(statement.indexName, statement.id, runId, lastRunStart, status, errMsg, rowCount)))
231239
if (!code.isSuccess()) {
232-
throw DataImportException("State manager failed, cannot log state for $statementId\n$response")
240+
throw DataImportException("State manager failed, cannot log state for ${statement.stateKey()}\n$response")
233241
}
234242
}
235243
}
Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,38 @@
1+
package uy.kohesive.elasticsearch.dataimport
2+
3+
import org.apache.spark.sql.SparkSession
4+
import org.apache.spark.sql.types.DataTypes
5+
import org.jsoup.Jsoup
6+
import org.jsoup.nodes.Document
7+
import org.jsoup.nodes.Entities
8+
import org.jsoup.parser.Parser
9+
import org.jsoup.safety.Whitelist
10+
import uy.kohesive.elasticsearch.dataimport.udf.Udfs
11+
12+
object DataImportHandlerUdfs {
13+
fun registerSparkUdfs(spark: SparkSession) {
14+
Udfs.registerStringToStringUdf(spark, "fluffly", fluffly)
15+
Udfs.registerStringToStringUdf(spark, "stripHtml", stripHtmlCompletely)
16+
Udfs.registerStringToStringUdf(spark, "unescapeHtmlEntites", unescapeHtmlEntities)
17+
}
18+
19+
@JvmStatic val whiteListMap = mapOf(
20+
"none" to Whitelist.none(),
21+
"basic" to Whitelist.basic(),
22+
"basicwithimages" to Whitelist.basicWithImages(),
23+
"relaxed" to Whitelist.relaxed(),
24+
"simpletext" to Whitelist.simpleText(),
25+
"simple" to Whitelist.simpleText()
26+
)
27+
28+
@JvmStatic val fluffly = fun (v: String): String = "fluffly " + v
29+
30+
@JvmStatic val stripHtmlCompletely = fun (v: String): String {
31+
return Jsoup.parseBodyFragment(v).text()
32+
}
33+
34+
@JvmStatic val unescapeHtmlEntities = fun (v: String): String {
35+
return Parser.unescapeEntities(v, false)
36+
}
37+
38+
}

0 commit comments

Comments
 (0)