Flink中使用基于名称的Row报错,如何正确传递至执行环境?
解决Flink中基于名称的Row传入执行环境的报错问题
错误原因
报错信息:
Caused by: java.lang.IllegalArgumentException: Accessing a field by position is not supported in name-based field mode.
直接使用env.fromCollection传入基于名称的Row时,Flink无法自动推断Row的字段结构(名称+类型),会默认尝试按位置访问字段,而名称模式的Row不支持这种访问方式,因此触发报错。
解决方案
必须给Flink明确指定基于名称的Row的Schema(即字段名和对应类型),通过创建RowTypeInfo并在fromCollection中传入即可解决问题,具体步骤如下:
1. 构建Row的Schema(RowTypeInfo)
从JsonNode中提取字段名,并匹配对应的Flink类型信息:
import org.apache.flink.api.common.typeinfo.BasicTypeInfo; import org.apache.flink.api.common.typeinfo.RowTypeInfo; import org.apache.flink.api.common.typeinfo.TypeInformation; import com.fasterxml.jackson.databind.JsonNode; import java.util.ArrayList; import java.util.List; // 收集字段名和对应类型信息 List<String> fieldNames = new ArrayList<>(); List<TypeInformation<?>> fieldTypes = new ArrayList<>(); jsonNode.fieldNames().forEachRemaining(fieldName -> { fieldNames.add(fieldName); JsonNode fieldNode = jsonNode.get(fieldName); // 根据Json节点类型匹配Flink内置类型,可按需扩展更多类型 if (fieldNode.isTextual()) { fieldTypes.add(BasicTypeInfo.STRING_TYPE_INFO); } else if (fieldNode.isInt()) { fieldTypes.add(BasicTypeInfo.INT_TYPE_INFO); } else if (fieldNode.isBoolean()) { fieldTypes.add(BasicTypeInfo.BOOLEAN_TYPE_INFO); } else { // 未匹配到的类型可使用通用OBJECT类型,或根据业务场景自定义处理 fieldTypes.add(BasicTypeInfo.OBJECT_TYPE_INFO); } }); // 创建携带字段名的RowTypeInfo RowTypeInfo rowTypeInfo = new RowTypeInfo( fieldTypes.toArray(new TypeInformation[0]), fieldNames.toArray(new String[0]) );
2. 传入带Schema的Row到执行环境
创建基于名称的Row后,在fromCollection方法中指定已定义好的rowTypeInfo:
Row row = Row.withNames(); jsonNode.fieldNames().forEachRemaining(fieldName -> { Object value; JsonNode fieldNode = jsonNode.get(fieldName); // 将Json节点值转换为对应Java类型 if (fieldNode.isTextual()) { value = fieldNode.asText(); } else if (fieldNode.isInt()) { value = fieldNode.asInt(); } else if (fieldNode.isBoolean()) { value = fieldNode.asBoolean(); } else { value = fieldNode; } row.setField(fieldName, value); }); StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 传入rowTypeInfo,让Flink识别基于名称的Row结构 env.fromCollection(List.of(row), rowTypeInfo);
内容的提问来源于stack exchange,提问作者ritratt
相关产品推荐
相关产品推荐

