schema只支持kafka自有域变量的问题,这个我修改的源码,思路是这样的:
从kafka中获取value域,将value通过from_json解析,然后将这个json作为Dataframe返回,不过,转为自定义的Dataframe时需要申明StructType,所以我又在配置文件中kafka部分增加了schema配置,用于记录StructType 的ddl。类似这样的格式:“source STRING, target STRING , rate DOUBLE”
目前这种方式跑通了,不过,我不知道这种方式是否合理。