apache/flink · error · ValidationException

Unexpected LogicalType:

Error message

Unexpected LogicalType: 

What it means

Validation rule: a non-repeated proto field of JavaType MESSAGE maps structurally to a Flink RowType. If the DDL declares any other type for such a field (scalar, array, map), this ValidationException is thrown with the actual LogicalType. It is raised at DDL validation time by PbSchemaValidationUtils.

Solutions

  1. Declare the column as ROW(field types...) matching the nested message's fields.
  2. Or use a computed column/dedicated view to flatten: inner_id AS inner.id if flattening is desired (map the ROW first).
  3. Regenerate and re-inspect the .proto to confirm the field is a message.

Example fix

-- before
inner VARCHAR,
-- after
inner ROW<street STRING, city STRING>,
Defensive patterns

Strategy: validation

Validate before calling

FieldDescriptor fd = descriptor.findFieldByName(name);
if (!fd.isRepeated() && fd.getJavaType() == FieldDescriptor.JavaType.MESSAGE
        && !(flinkType instanceof RowType)) {
    throw new ValidationException(name + " must be ROW");
}

Type guard

boolean matchesProtoShape(FieldDescriptor fd, LogicalType t) {
    if (!fd.isRepeated() && fd.getJavaType() == FieldDescriptor.JavaType.MESSAGE) {
        return t instanceof RowType;
    }
    return true;
}

Prevention

When it happens

Trigger: Proto: message Outer { Inner inner = 1; } with DDL column inner INT/VARCHAR/ARRAY<...> instead of ROW(...).

Common situations: Hand-written DDLs that flatten message fields as scalars; schema drift after a proto field changed from scalar to message.

Related errors


AI-assisted analysis of apache/flink@2f3c205e92 (2026-08-14). Data as JSON: /api/errors/130704c0a36efa42. Report an issue: GitHub.

Appendix: source

Thrown at flink-formats/flink-protobuf/src/main/java/org/apache/flink/formats/protobuf/util/PbSchemaValidationUtils.java:105

                            }
                        });
    }

    /**
     * Validate type match of general type.
     *
     * @param fd the {@link Descriptors.Descriptor} of the protobuf object.
     * @param logicalType the corresponding {@link LogicalType} to the {@link FieldDescriptor}
     */
    private static void validateTypeMatch(FieldDescriptor fd, LogicalType logicalType) {
        if (!fd.isRepeated()) {
            if (fd.getJavaType() != JavaType.MESSAGE) {
                // simple type
                validateSimpleType(fd, logicalType.getTypeRoot());
            } else {
                // message type
                if (!(logicalType instanceof RowType)) {
                    throw new ValidationException(
                            "Unexpected LogicalType: " + logicalType + ". It should be RowType");
                }
                validateTypeMatch(fd.getMessageType(), (RowType) logicalType);
            }
        } else {
            if (fd.isMapField()) {
                // map type
                if (!(logicalType instanceof MapType)) {
                    throw new ValidationException(
                            "Unexpected LogicalType: " + logicalType + ". It should be MapType");
                }
                MapType mapType = (MapType) logicalType;
                validateSimpleType(
                        fd.getMessageType().findFieldByName(PbConstant.PB_MAP_KEY_NAME),
                        mapType.getKeyType().getTypeRoot());
                validateTypeMatch(
                        fd.getMessageType().findFieldByName(PbConstant.PB_MAP_VALUE_NAME),
                        mapType.getValueType());

View on GitHub (pinned to 2f3c205e92)