Amyth
08/08/2024, 5:21 AMimport org.apache.flink.api.common.typeinfo.TypeInformation
import org.apache.flink.api.scala.typeutils.Types
import org.apache.flink.streaming.api.scala.DataStream
import org.apache.flink.table.api.Table
import org.apache.flink.table.api.bridge.scala.StreamTableEnvironment
import org.apache.flink.types.Row
import scala.annotation.tailrec
object Functions extends Serializable {
def convertDatastreamToTable(datastream: DataStream[_], tableEnv: StreamTableEnvironment, columnNames: String): Table = { tableEnv.fromDataStream(datastream).as(columnNames)
}
tableEnv.fromDataStream(datastream).as(columnNames).getResolvedSchema
prints-
( timestamp TIMESTAMP(9),
unit FLOAT,
name STRING,
ip STRING,
`c_id`STRING,
app_name STRING,
id STRING,
host STRING )
but the Flink version upgrade to Flink-1.18.1
import org.apache.flink.table.api.*
import org.apache.flink.api.common.typeinfo.{TypeInformation, Types}
import org.apache.flink.streaming.api.datastream.DataStream
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment
import org.apache.flink.table.types.DataType
import org.apache.flink.types.Row import scala.annotation.tailrec
object Functions extends Serializable {
def convertDatastreamToTable(datastream: DataStream[_], tableEnv: StreamTableEnvironment, columnNames: String): Table = { tableEnv.fromDataStream(datastream).as(columnNames)
}
tableEnv.fromDataStream(datastream).as(columnNames).getResolvedSchema
I am getting-
( timestamp,unit, name,ip,c_id,app_name, id,host, RAW('org.apache.flink.types.Row', '...') )
I want my output to be like 1.14.0,
In 1.14.0 I was using scala api's but in Flink 1.18.1 I am using Java APIs. Please suggest