美文网首页Spark在简书大数据
spark structured streaming的sourc

spark structured streaming的sourc

作者: 中科院_白乔 | 来源:发表于2017-09-05 20:22 被阅读33次

下面是一段创建structured streaming的Dataset的代码:

val lines = spark.readStream.format("socket")
    .option("host", "localhost").option("port", 9999).load();

会创建一个socket类型的Source,该name2class的映射由DataSource.lookupDataSource()完成

val serviceLoader = ServiceLoader.load(classOf[DataSourceRegister], loader)
...
serviceLoader.asScala.filter(_.shortName().equalsIgnoreCase(provider1)).toList
...

应该是从当前类路径中查找所有的DataSourceRegister,并读取它的shortName(),如果是"socket"就确定了由该DataSourceRegister来创建对应的DataSource

果然,有一个TextSocketSourceProvider

class TextSocketSourceProvider extends StreamSourceProvider with DataSourceRegister with Logging {
...
override def shortName(): String = "socket"

  override def createSource(
      sqlContext: SQLContext,
      metadataPath: String,
      schema: Option[StructType],
      providerName: String,
      parameters: Map[String, String]): Source = {
    val host = parameters("host")
    val port = parameters("port").toInt
    new TextSocketSource(host, port, parseIncludeTimestamp(parameters), sqlContext)
  }
}

TextSocketSourceProvider的createSource创建一个TextSocketSource

TextSocketSource是一个Source,Source接口如下:

trait Source  {
  def schema: StructType
  def getOffset: Option[Offset]
  def getBatch(start: Option[Offset], end: Offset): DataFrame
  def commit(end: Offset) : Unit = {}
  def stop(): Unit
}

相关文章

网友评论

    本文标题:spark structured streaming的sourc

    本文链接:https://www.haomeiwen.com/subject/rgspjxtx.html