Skip to content

Commit

Permalink
feat(jdbc): add postQuery parameter (#414)
Browse files Browse the repository at this point in the history
  • Loading branch information
lyogev committed Feb 28, 2021
1 parent 59d2bd1 commit 26969b0
Show file tree
Hide file tree
Showing 2 changed files with 15 additions and 1 deletion.
2 changes: 2 additions & 0 deletions config/metric_config_sample.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,8 @@ output:
outputOptions:
saveMode: Overwrite
dbTable: table_in_db
# Optional: run this after writing to JDBC
postQuery: CREATE INDEX...
- dataFrameName: table1
outputType: JDBCQuery
outputOptions:
Expand Down
Original file line number Diff line number Diff line change
@@ -1,12 +1,13 @@
package com.yotpo.metorikku.output.writers.jdbc

import java.util.Properties

import com.yotpo.metorikku.configuration.job.output.JDBC
import com.yotpo.metorikku.output.Writer
import org.apache.log4j.LogManager
import org.apache.spark.sql.{DataFrame, SaveMode}

import java.sql.DriverManager


class JDBCOutputWriter(props: Map[String, String], jdbcConf: Option[JDBC]) extends Writer {

Expand All @@ -26,6 +27,17 @@ class JDBCOutputWriter(props: Map[String, String], jdbcConf: Option[JDBC]) exten
val writer = df.write.format(jdbcConf.driver)
.mode(dbOptions.saveMode)
.jdbc(jdbcConf.connectionUrl, dbOptions.dbTable, connectionProperties)

props.get("postQuery") match {
case Some(query) =>
val conn = DriverManager.getConnection(jdbcConf.connectionUrl, jdbcConf.user, jdbcConf.password)
val stmt = conn.prepareStatement(query)
stmt.execute()
stmt.close()
conn.close()
case _ =>
}

case None =>
}
}
Expand Down

0 comments on commit 26969b0

Please sign in to comment.