【发布时间】:2020-07-17 20:37:39
【问题描述】:
我的任务是从 jdbc 中获取表,然后将它们放到 s3 中。我已经使用 Slick 模式代码生成器生成了这些表的类。 如果我为每个表手动编写代码,它就可以完美运行。 就这样
Slick
.source(Tables.table1.result)
.runWith(ParquetStreams.toParquetSingleFile(s"s3a://bucket/table1"))
.onComplete {
case _ =>
println("table1")
}
Slick
.source(Tables.table2.result)
.runWith(ParquetStreams.toParquetSingleFile(s"s3a://bucket/table2"))
.onComplete {
case _ =>
println("table2")
}
问题是我有很多表,如果我可以迭代它们会更容易。
代码
val tbls = Map("table1" -> Tables.table1, "table2" -> Tables.table2)
tbls.foreach(table => {
val table_name = table._1
Slick
.source(table_2.result)
.runWith(ParquetStreams.toParquetSingleFile(s"s3a://bucket/$table_name"))
.onComplete {
case _ =>
println(table_name)
}
})
收到此错误。
could not find implicit value for evidence parameter of type com.github.mjakubowski84.parquet4s.ParquetRecordEncoder[_1#TableElementType]
[error] .runWith(ParquetStreams.toParquetSingleFile(s"s3a://bucket/$table_name"))
编辑:------------------------------------------ ------------------
感谢您澄清事情。但即使有建议也行不通。应用这些东西后,存在同样的错误,并且出现了更多错误。
found : slick.lifted.TableQuery[_1] where type _1 >: Tables.TableTwo with Tables.TableOne <: Tables.profile.Table[_ >: Tables.TableTwoRow with Tables.TableOneRow <: Product with java.io.Serializable]
按照某人的要求在此处留下简化代码。
import java.sql.Timestamp
import akka.actor.typed.ActorSystem
import akka.actor.typed.scaladsl.Behaviors
import akka.stream.alpakka.slick.javadsl.SlickSession
import akka.stream.alpakka.slick.scaladsl.Slick
import scala.concurrent.duration._
import scala.concurrent.{Await, ExecutionContext}
import com.github.mjakubowski84.parquet4s.{ParquetRecordEncoder, ParquetSchemaResolver, ParquetStreams}
object Tables extends {
val profile = slick.jdbc.MySQLProfile
} with Tables
/** Slick data model trait for extension, choice of backend or usage in the cake pattern. (Make sure to initialize this late.) */
trait Tables {
val profile: slick.jdbc.JdbcProfile
import profile.api._
case class TableOneRow(id: Int, values: Option[String] = None)
class TableOne(_tableTag: Tag) extends profile.api.Table[TableOneRow](_tableTag, Some("schema"), "table_one") {
def * = (id, values) <> (TableOneRow.tupled, TableOneRow.unapply)
val id: Rep[Int] = column[Int]("id", O.AutoInc, O.PrimaryKey)
val values: Rep[Option[String]] = column[Option[String]]("values", O.Default(None))
}
lazy val TableOne = new TableQuery(tag => new TableOne(tag))
case class TableTwoRow(id: Int, values: Option[Timestamp] = None)
class TableTwo(_tableTag: Tag) extends profile.api.Table[TableTwoRow](_tableTag, Some("schema"), "table_two") {
def * = (id, date) <> (TableTwoRow.tupled, TableTwoRow.unapply)
val id: Rep[Int] = column[Int]("id", O.AutoInc, O.PrimaryKey)
val date: Rep[Option[Timestamp]] = column[Option[Timestamp]]("date", O.Default(None))
}
lazy val TableTwo = new TableQuery(tag => new TableTwo(tag))
}
object Main extends App {
implicit val actorSystem: ActorSystem[Nothing] = ActorSystem(Behaviors.empty, "alpakka-sample")
implicit val executionContext: ExecutionContext = actorSystem.executionContext
implicit val session = SlickSession.forConfig("slick-mysql") // (1)
import session.profile.api._
case class TableWithRecordEncoder[A](
table: TableQuery[A])(
implicit val recordEncoder: ParquetRecordEncoder[A]
)
def doTheThings[A](table: TableWithRecordEncoder[A], path: String) = {
import table.recordEncoder
Slick
.source(table.table.result)
.runWith(ParquetStreams.toParquetSingleFile(path))
.onComplete {
case _ =>
println("Done. " + path)
}
}
import polymorphic._
def withRecordEncoder[A](table: TableQuery[A])(implicit recordEncoder: ParquetRecordEncoder[A])
: Exists[TableWithRecordEncoder]
= Exists(TableWithRecordEncoder(table))
import polymorphic.syntax.all._
Map("s3a://bucket/table_1" -> withRecordEncoder(Tables.TableOne),
"s3a://bucket/table_2" -> withRecordEncoder(Tables.TableTwo)).
foreach{ case (path, table) =>
doTheThings(table.value, path)
}
actorSystem.whenTerminated.map(_ => session.close())
Await.result(actorSystem.whenTerminated, Duration.Inf)
}
【问题讨论】:
-
你能发布完整的源代码文件吗?如果第一个示例编译,则该示例中的范围内必须有一个隐式值,该值在重构后不再在范围内。如果你提取了一个方法,你可能需要给它添加一个隐式参数。
-
再看,我认为这是推断类型的问题,因为当您将它们添加到映射时,它会丢失每个表值的具体类型并扩大到超类型。那么很可能没有为该超类型定义通用编码器。
标签: scala amazon-s3 akka parquet alpakka