restapi(8)- restapi-sql:用户自主的服务

2024-04-09 04:38

本文主要是介绍restapi(8)- restapi-sql:用户自主的服务,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

  学习函数式编程初衷是看到自己熟悉的oop编程语言和sql数据库在现代商业社会中前景暗淡,准备完全放弃windows技术栈转到分布式大数据技术领域的,但是在现实中理想总是不如人意的。本来想在一个规模较小的公司展展拳脚,以为小公司会少点历史包袱,有利于全面技术改造。但现实是:即使是小公司,一旦有个成熟的产品,那么进行全面的技术更新基本上是不可能的了,因为公司要生存,开发人员很难新旧技术之间随时切换。除非有狂热的热情,员工怠慢甚至抵制情绪不容易解决。只能采取逐步切换方式:保留原有产品的后期维护不动,新产品开发用一些新的技术。在我们这里的情况就是:以前一堆c#、sqlserver的东西必须保留,新的功能比如大数据、ai、识别等必须用新的手段如scala、python、dart、akka、kafka、cassandra、mongodb来开发。好了,新旧两个开发平台之间的软件系统对接又变成了一个问题。

   现在我们这里又个需求:把在linux-ubuntu akka-cluster集群环境里mongodb里数据处理的结果传给windows server下SQLServer里。这是一种典型的异系统集成场景。我的解决方案是通过一个restapi服务作为两个系统的数据桥梁,这个restapi的最基本要求是:

1、支持任何操作系统前端:这个没什么问题,在http层上通过json交换数据

2、能读写mongodb:在前面讨论的restapi-mongo已经实现了这一功能

3、能读写windows server环境下的sqlserver:这个是本篇讨论的主题

前面曾经实现了一个jdbc-engine项目,基于scalikejdbc,不过只示范了slick-h2相关的功能。现在需要sqlserver-jdbc驱动,然后试试能不能在JVM里驱动windows下的sqlserver。maven里找不到sqlserver的驱动,但从微软官网可以下载mssql-jdbc-7.0.0.jre8.jar。这是个jar,在sbt里称作unmanagedjar,不能摆在build.sbt的dependency里。这个需要摆在项目根目录下的lib目录下即可(也可以在放在build.sbt里unmanagedBase :=?? 指定的路径下)。然后是数据库连接,下面是可以使用sqlserver的application.conf配置文件内容:

# JDBC settings
prod {db {h2 {driver = "org.h2.Driver"url = "jdbc:h2:tcp://localhost/~/slickdemo"user = ""password = ""poolFactoryName = "hikaricp"numThreads = 10maxConnections = 12minConnections = 4keepAliveConnection = true}mysql {driver = "com.mysql.cj.jdbc.Driver"url = "jdbc:mysql://localhost:3306/testdb"user = "root"password = "123"poolFactoryName = "hikaricp"numThreads = 10maxConnections = 12minConnections = 4keepAliveConnection = true}postgres {driver = "org.postgresql.Driver"url = "jdbc:postgresql://localhost:5432/testdb"user = "root"password = "123"poolFactoryName = "hikaricp"numThreads = 10maxConnections = 12minConnections = 4keepAliveConnection = true}mssql {driver = "com.microsoft.sqlserver.jdbc.SQLServerDriver"url = "jdbc:sqlserver://192.168.11.164:1433;integratedSecurity=false;Connect Timeout=3000"user = "sa"password = "Tiger2020"poolFactoryName = "hikaricp"numThreads = 10maxConnections = 12minConnections = 4keepAliveConnection = trueconnectionTimeout = 3000}termtxns {driver = "com.microsoft.sqlserver.jdbc.SQLServerDriver"url = "jdbc:sqlserver://192.168.11.164:1433;DATABASE=TERMTXNS;integratedSecurity=false;Connect Timeout=3000"user = "sa"password = "Tiger2020"poolFactoryName = "hikaricp"numThreads = 10maxConnections = 12minConnections = 4keepAliveConnection = trueconnectionTimeout = 3000}crmdb {driver = "com.microsoft.sqlserver.jdbc.SQLServerDriver"url = "jdbc:sqlserver://192.168.11.164:1433;DATABASE=CRMDB;integratedSecurity=false;Connect Timeout=3000"user = "sa"password = "Tiger2020"poolFactoryName = "hikaricp"numThreads = 10maxConnections = 12minConnections = 4keepAliveConnection = trueconnectionTimeout = 3000}}# scallikejdbc Global settingsscalikejdbc.global.loggingSQLAndTime.enabled = truescalikejdbc.global.loggingSQLAndTime.logLevel = infoscalikejdbc.global.loggingSQLAndTime.warningEnabled = truescalikejdbc.global.loggingSQLAndTime.warningThresholdMillis = 1000scalikejdbc.global.loggingSQLAndTime.warningLogLevel = warnscalikejdbc.global.loggingSQLAndTime.singleLineMode = falsescalikejdbc.global.loggingSQLAndTime.printUnprocessedStackTrace = falsescalikejdbc.global.loggingSQLAndTime.stackTraceDepth = 10
}

这个文件里的mssql,termtxns,crmdb段落都是给sqlserver的,它们都使用hikaricp线程池管理。

在jdbc-engine里启动数据库方式如下:

  ConfigDBsWithEnv("prod").setup('termtxns)ConfigDBsWithEnv("prod").setup('crmdb)ConfigDBsWithEnv("prod").loadGlobalSettings()

这段打开了在配置文件中用termtxns,crmdb注明的数据库。

下面是SqlHttpServer.scala的代码:

package com.datatech.rest.sql
import akka.http.scaladsl.Http
import akka.http.scaladsl.server.Directives._
import pdi.jwt._
import AuthBase._
import MockUserAuthService._
import com.datatech.sdp.jdbc.config.ConfigDBsWithEnvimport akka.actor.ActorSystem
import akka.stream.ActorMaterializerimport Repo._
import SqlRoute._object SqlHttpServer extends App {implicit val httpSys = ActorSystem("sql-http-sys")implicit val httpMat = ActorMaterializer()implicit val httpEC = httpSys.dispatcherConfigDBsWithEnv("prod").setup('termtxns)ConfigDBsWithEnv("prod").setup('crmdb)ConfigDBsWithEnv("prod").loadGlobalSettings()implicit val authenticator = new AuthBase().withAlgorithm(JwtAlgorithm.HS256).withSecretKey("OpenSesame").withUserFunc(getValidUser)val route =path("auth") {authenticateBasic(realm = "auth", authenticator.getUserInfo) { userinfo =>post { complete(authenticator.issueJwt(userinfo))}}} ~pathPrefix("api") {authenticateOAuth2(realm = "api", authenticator.authenticateToken) { token =>new SqlRoute("sql", token)(new JDBCRepo).route// ~ ...}}val (port, host) = (50081,"192.168.11.189")val bindingFuture = Http().bindAndHandle(route,host,port)println(s"Server running at $host $port. Press any key to exit ...")scala.io.StdIn.readLine()bindingFuture.flatMap(_.unbind()).onComplete(_ => httpSys.terminate())}

服务入口在http://mydemo.com/api/sql,服务包括get,post,put三类,看看这个SqlRoute:

package com.datatech.rest.sql
import akka.http.scaladsl.server.Directives
import akka.stream.ActorMaterializer
import akka.http.scaladsl.model._
import akka.actor.ActorSystem
import com.datatech.rest.sql.Repo.JDBCRepo
import akka.http.scaladsl.common._
import spray.json.DefaultJsonProtocol
import akka.http.scaladsl.marshallers.sprayjson.SprayJsonSupporttrait JsFormats extends SprayJsonSupport with DefaultJsonProtocol
object JsConverters extends JsFormats {import SqlModels._implicit val brandFormat = jsonFormat2(Brand)implicit val customerFormat = jsonFormat6(Customer)
}object SqlRoute {import JsConverters._implicit val jsonStreamingSupport = EntityStreamingSupport.json().withParallelMarshalling(parallelism = 8, unordered = false)class SqlRoute(val pathName: String, val jwt: String)(repo: JDBCRepo)(implicit  sys: ActorSystem, mat: ActorMaterializer) extends Directives with JsonConverter {val route = pathPrefix(pathName) {path(Segment / Remaining) { case (db, tbl) =>(get & parameter('sqltext)) { sql => {val rsc = new RSConverterval rows = repo.query[Map[String,Any]](db, sql, rsc.resultSet2Map)complete(rows.map(m => toJson(m)))}} ~ (post & parameter('sqltext)) { sql =>entity(as[String]){ json =>repo.batchInsert(db,tbl,sql,json)complete(StatusCodes.OK)}} ~ put {entity(as[Seq[String]]) { sqls =>repo.update(db, sqls)complete(StatusCodes.OK)}}}}}
}

jdbc-engine的特点是可以用字符类型的sql语句来操作。所以我们可以通过传递字符串型的sql语句来实现服务调用,很通用。restapi-sql提供的是对服务器端sqlserver的普通操作,包括读get,写入post,更改put。这些sqlserver操作部分是在JDBCRepo里的:

package com.datatech.rest.sql
import com.datatech.sdp.jdbc.engine.JDBCEngine._
import com.datatech.sdp.jdbc.engine.{JDBCQueryContext, JDBCUpdateContext}
import scalikejdbc._
import akka.stream.ActorMaterializer
import com.datatech.sdp.result.DBOResult.DBOResult
import akka.stream.scaladsl._
import scala.concurrent._
import SqlModels._object Repo {class JDBCRepo(implicit ec: ExecutionContextExecutor, mat: ActorMaterializer) {def query[R](db: String, sqlText: String, toRow: WrappedResultSet => R): Source[R,Any] = {//construct the contextval ctx = JDBCQueryContext(dbName = Symbol(db),statement = sqlText)jdbcAkkaStream(ctx,toRow)}def query(db: String, tbl: String, sqlText: String) = {//construct the contextval ctx = JDBCQueryContext(dbName = Symbol(db),statement = sqlText)jdbcQueryResult[Vector,RS](ctx,getConverter(tbl)).toFuture[Vector[RS]]}def update(db: String, sqlTexts: Seq[String]): DBOResult[Seq[Long]] = {val ctx = JDBCUpdateContext(dbName = Symbol(db),statements = sqlTexts)jdbcTxUpdates(ctx)}def bulkInsert[P](db: String, sqlText: String, prepParams: P => Seq[Any], params: Source[P,_]) = {val insertAction = JDBCActionStream(dbName = Symbol(db),parallelism = 4,processInOrder = false,statement = sqlText,prepareParams = prepParams)params.via(insertAction.performOnRow).to(Sink.ignore).run()}def batchInsert(db: String, tbl: String, sqlText: String, jsonParams: String):DBOResult[Seq[Long]] = {val ctx = JDBCUpdateContext(dbName = Symbol(db),statements = Seq(sqlText),batch = true,parameters = getSeqParams(jsonParams,sqlText))jdbcBatchUpdate[Seq](ctx)}}import monix.execution.Scheduler.Implicits.globalimplicit class DBResultToFuture(dbr: DBOResult[_]){def toFuture[R] = {dbr.value.value.runToFuture.map {eor =>eor match {case Right(or) => or match {case Some(r) => r.asInstanceOf[R]case None => throw new RuntimeException("Operation produced None result!")}case Left(err) => throw new RuntimeException(err)}}}}
}

读query部分即 def query[R](db: String, sqlText: String, toRow: WrappedResultSet => R): Source[R,Any] = {...} 这个函数返回Source[R,Any],下面我们好好谈谈这个R:R是读的结果,通常是某个类或model,比如读取Person记录返回一组Person类的实例。这里有一种强类型的感觉。一开始我也是随大流坚持建model后用toJson[E],fromJson[E]这样做线上数据转换。现在的问题是restapi-sql是一项公共服务,使用者知道sqlserver上有些什么表,然后希望通过sql语句来从这些表里读取数据。这些sql语句可能超出表的界限如sql join, union等,如果我们坚持每个返回结果都必须有个对应的model,那么显然就会牺牲这个服务的通用性。实际上,http线上数据交换本身就不可能是强类型的,因为经过了json转换。对于json转换来说,只要求字段名称、字段类型对称就行了。至于从什么类型转换成了另一个什么类型都没问题。所以,字段名+字段值的表现形式不就是Map[K,V]吗,我们就用Map[K,V]作为万能model就行了,没人知道。也就是说我们可以把jdbc的ResultSet转成Map[K,V]然后再转成json,接收方可以获取与model同样的字段名和字段值。好,就把ResultSet转成Map[String,Any]:

package com.datatech.rest.sql
import scalikejdbc._
import java.sql.ResultSetMetaData
class RSConverter {import RSConverterUtil._var rsMeta: ResultSetMetaData = _var columnCount: Int = 0var rsFields: List[(String,String)] = List[(String,String)]()def getFieldsInfo:List[(String,String)] =( 1 until columnCount).foldLeft(List[(String,String)]()) {case (cons,i) =>(rsMeta.getColumnName(i) -> rsMeta.getColumnTypeName(i)) :: cons}def resultSet2Map(rs: WrappedResultSet): Map[String,Any] = {if(columnCount == 0) {rsMeta =  rs.underlying.getMetaDatacolumnCount = rsMeta.getColumnCountrsFields = getFieldsInfo}rsFields.foldLeft(Map[String,Any]()) {case (m,(n,t)) =>m + (n -> rsFieldValue(n,t,rs))}}
}
object RSConverterUtil {import scala.collection.immutable.TreeMapdef map2Params(stm: String, m: Map[String,Any]): Seq[Any] = {val sortedParams = m.foldLeft(TreeMap[Int,Any]()) {case (t,(k,v)) => t + (stm.indexOfSlice(k) -> v)}sortedParams.map(_._2).toSeq}def rsFieldValue(fldname: String, fldType: String, rs: WrappedResultSet): Any = fldType match {case "LONGVARCHAR" => rs.string(fldname)case "VARCHAR" => rs.string(fldname)case "CHAR" => rs.string(fldname)case "BIT" => rs.boolean(fldname)case "TIME" => rs.time(fldname)case "TIMESTAMP" => rs.timestamp(fldname)case "ARRAY" => rs.array(fldname)case "NUMERIC" => rs.bigDecimal(fldname)case "BLOB" => rs.blob(fldname)case "TINYINT" => rs.byte(fldname)case "VARBINARY" => rs.bytes(fldname)case "BINARY" => rs.bytes(fldname)case "CLOB" => rs.clob(fldname)case "DATE" => rs.date(fldname)case "DOUBLE" => rs.double(fldname)case "REAL" => rs.float(fldname)case "FLOAT" => rs.float(fldname)case "INTEGER" => rs.int(fldname)case "SMALLINT" => rs.int(fldname)case "Option[Int]" => rs.intOpt(fldname)case "BIGINT" => rs.long(fldname)}
}

下面是个调用query服务的例子:

    val getAllRequest = HttpRequest(HttpMethods.GET,uri = "http://192.168.11.189:50081/api/sql/termtxns/brand?sqltext=SELECT%20*%20FROM%20BRAND",).addHeader(authentication)(for {response <- Http().singleRequest(getAllRequest)message <- Unmarshal(response.entity).to[String]} yield message).andThen {case Success(msg) => println(s"Received message: $msg")case Failure(err) => println(s"Error: ${err.getMessage}")}

特点是我只需要提供sql语句,服务就会返回一个json数组,然后我怎么转换json就随我高兴了。

这篇关于restapi(8)- restapi-sql:用户自主的服务的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/887194

相关文章

SQL中的外键约束

外键约束用于表示两张表中的指标连接关系。外键约束的作用主要有以下三点: 1.确保子表中的某个字段(外键)只能引用父表中的有效记录2.主表中的列被删除时,子表中的关联列也会被删除3.主表中的列更新时,子表中的关联元素也会被更新 子表中的元素指向主表 以下是一个外键约束的实例展示

基于MySQL Binlog的Elasticsearch数据同步实践

一、为什么要做 随着马蜂窝的逐渐发展,我们的业务数据越来越多,单纯使用 MySQL 已经不能满足我们的数据查询需求,例如对于商品、订单等数据的多维度检索。 使用 Elasticsearch 存储业务数据可以很好的解决我们业务中的搜索需求。而数据进行异构存储后,随之而来的就是数据同步的问题。 二、现有方法及问题 对于数据同步,我们目前的解决方案是建立数据中间表。把需要检索的业务数据,统一放到一张M

如何去写一手好SQL

MySQL性能 最大数据量 抛开数据量和并发数,谈性能都是耍流氓。MySQL没有限制单表最大记录数,它取决于操作系统对文件大小的限制。 《阿里巴巴Java开发手册》提出单表行数超过500万行或者单表容量超过2GB,才推荐分库分表。性能由综合因素决定,抛开业务复杂度,影响程度依次是硬件配置、MySQL配置、数据表设计、索引优化。500万这个值仅供参考,并非铁律。 博主曾经操作过超过4亿行数据

性能分析之MySQL索引实战案例

文章目录 一、前言二、准备三、MySQL索引优化四、MySQL 索引知识回顾五、总结 一、前言 在上一讲性能工具之 JProfiler 简单登录案例分析实战中已经发现SQL没有建立索引问题,本文将一起从代码层去分析为什么没有建立索引? 开源ERP项目地址:https://gitee.com/jishenghua/JSH_ERP 二、准备 打开IDEA找到登录请求资源路径位置

MySQL数据库宕机,启动不起来,教你一招搞定!

作者介绍:老苏,10余年DBA工作运维经验,擅长Oracle、MySQL、PG、Mongodb数据库运维(如安装迁移,性能优化、故障应急处理等)公众号:老苏畅谈运维欢迎关注本人公众号,更多精彩与您分享。 MySQL数据库宕机,数据页损坏问题,启动不起来,该如何排查和解决,本文将为你说明具体的排查过程。 查看MySQL error日志 查看 MySQL error日志,排查哪个表(表空间

MySQL高性能优化规范

前言:      笔者最近上班途中突然想丰富下自己的数据库优化技能。于是在查阅了多篇文章后,总结出了这篇! 数据库命令规范 所有数据库对象名称必须使用小写字母并用下划线分割 所有数据库对象名称禁止使用mysql保留关键字(如果表名中包含关键字查询时,需要将其用单引号括起来) 数据库对象的命名要能做到见名识意,并且最后不要超过32个字符 临时库表必须以tmp_为前缀并以日期为后缀,备份

【区块链 + 人才服务】可信教育区块链治理系统 | FISCO BCOS应用案例

伴随着区块链技术的不断完善,其在教育信息化中的应用也在持续发展。利用区块链数据共识、不可篡改的特性, 将与教育相关的数据要素在区块链上进行存证确权,在确保数据可信的前提下,促进教育的公平、透明、开放,为教育教学质量提升赋能,实现教育数据的安全共享、高等教育体系的智慧治理。 可信教育区块链治理系统的顶层治理架构由教育部、高校、企业、学生等多方角色共同参与建设、维护,支撑教育资源共享、教学质量评估、

[MySQL表的增删改查-进阶]

🌈个人主页:努力学编程’ ⛅个人推荐: c语言从初阶到进阶 JavaEE详解 数据结构 ⚡学好数据结构,刷题刻不容缓:点击一起刷题 🌙心灵鸡汤:总有人要赢,为什么不能是我呢 💻💻💻数据库约束 🔭🔭🔭约束类型 not null: 指示某列不能存储 NULL 值unique: 保证某列的每行必须有唯一的值default: 规定没有给列赋值时的默认值.primary key:

【区块链 + 人才服务】区块链集成开发平台 | FISCO BCOS应用案例

随着区块链技术的快速发展,越来越多的企业开始将其应用于实际业务中。然而,区块链技术的专业性使得其集成开发成为一项挑战。针对此,广东中创智慧科技有限公司基于国产开源联盟链 FISCO BCOS 推出了区块链集成开发平台。该平台基于区块链技术,提供一套全面的区块链开发工具和开发环境,支持开发者快速开发和部署区块链应用。此外,该平台还可以提供一套全面的区块链开发教程和文档,帮助开发者快速上手区块链开发。

MySQL-CRUD入门1

文章目录 认识配置文件client节点mysql节点mysqld节点 数据的添加(Create)添加一行数据添加多行数据两种添加数据的效率对比 数据的查询(Retrieve)全列查询指定列查询查询中带有表达式关于字面量关于as重命名 临时表引入distinct去重order by 排序关于NULL 认识配置文件 在我们的MySQL服务安装好了之后, 会有一个配置文件, 也就