apache / apache/iotdb

read empty dataframe when one column result by spark iotdb connector

Offen
#2,861 3 Kommentare 0 Reaktionen 0 zugewiesene Personen Auf GitHub ansehen
Module - Spark
Vorherrschende Sprache
Java
Sterne
6.4k
Forks
1.2k
Ø Merge
1 T. 23 Std.
Gemergte PRs (30 T.)
115

Beschreibung

**环境**
- Iotdb 0.9.3 & 0.11.2
- spark 2.3.3
- spark-iotdb-connector 0.11.2

**问题**
spark读取iotdb数据后,返回一个空的DataFrame对象

**分析**
0.9.3中执行带聚合函数的sql,如`select count(root.ln.wf01.wt01.status) from root`,返回为两列,一列是时间列,一列是值,但是在0.11.2中返回的是一列,但是spark-iotdb-connector中并未处理这种,在生成schame的时候,从2开始遍历,导致只有一列的数据返回空schema,代码如下:
spark-iotdb-connector\src\main\scala\org\apache\iotdb\spark\db\Converter.scala
```
def toSparkSchema(options: IoTDBOptions): StructType = {

Class.forName("org.apache.iotdb.jdbc.IoTDBDriver")
val sqlConn: Connection = DriverManager.getConnection(options.url, options.user, options.password)
val sqlStatement: Statement = sqlConn.createStatement()
val hasResultSet: Boolean = sqlStatement.execute(options.sql)

val fields = new ListBuffer[StructField]()
if (hasResultSet) {
val resultSet: ResultSet = sqlStatement.getResultSet
val resultSetMetaData: ResultSetMetaData = resultSet.getMetaData

val printTimestamp = !resultSet.asInstanceOf[IoTDBJDBCResultSet].isIgnoreTimeStamp
if (printTimestamp) {
fields += StructField(SQLConstant.TIMESTAMP_STR, LongType, nullable = false)
}

val colCount = resultSetMetaData.getColumnCount
for (i <- 2 to colCount) {
.....
}
StructType(fields.toList)
}
else {
StructType(fields)
}
}
```
需要加个startCol变量,来确定是从2开始遍历还是1,如下:
```
var startCol = 1
if (printTimestamp) {
fields += StructField(SQLConstant.TIMESTAMP_STR, LongType, nullable = false)
startCol = 2
}

val colCount = resultSetMetaData.getColumnCount
for (i <- startCol to colCount) {
.....
}
```

Beitragsleitfaden

Beitragsleitfaden öffnen

Rechercherichtung

Beginne mit spark-iotdb-connector/src/main/scala/org/apache/iotdb/spark/db/Converter.scala und untersuche toSparkSchema. Führe dann die gemeldete Aggregatabfrage mit den aufgeführten IoTDB- und Spark-Versionen erneut aus. Überprüfe, dass ein Ergebnis mit einer Spalte ein nichtleeres Spark-Schema und einen DataFrame anstelle eines leeren Schemas erzeugt.

Vom Indexierungsmodell aus dem Issue-Text verfasst.

Bewertung

Tech-Stack
scala, spark, sql
Bereich
databases
Issue-Typ
Bug
Schwierigkeit
2/5
Geschätzter Aufwand
1-3 Stunden
Aktivitätsstatus
Veraltet
Klarheit
Klar beschrieben
Anfängerfreundlichkeit
58/100

Neue Issues direkt in Ihr Postfach

Eine kurze Übersicht über anfängerfreundliche GitHub-Issues.