Github user HyukjinKwon commented on a diff in the pull request:
https://github.com/apache/spark/pull/20397#discussion_r167393259
--- Diff:
sql/core/src/main/scala/org/apache/spark/sql/execution/streaming/sources/RateStreamSourceV2.scala
---
@@ -139,15 +139,15 @@ class RateStreamMicroBatchReader(options:
DataSourceV2Options)
outTimeMs += msPerPartitionBetweenRows
}
- RateStreamBatchTask(packedRows).asInstanceOf[ReadTask[Row]]
+ RateStreamBatchTask(packedRows).asInstanceOf[DataReaderFactory[Row]]
}.toList.asJava
}
override def commit(end: Offset): Unit = {}
override def stop(): Unit = {}
}
-case class RateStreamBatchTask(vals: Seq[(Long, Long)]) extends
ReadTask[Row] {
+case class RateStreamBatchTask(vals: Seq[(Long, Long)]) extends
DataReaderFactory[Row] {
--- End diff --
and should we rename `RateStreamBatchTask` too?
---
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]