import org.apache.flink.core.io.SimpleVersionedSerializer;
import org.apache.flink.streaming.api.functions.sink.filesystem.BucketAssigner;
import org.apache.flink.streaming.api.functions.sink.filesystem.bucketassigners.SimpleVersionedStringSerializer;
public class FileBucketAssigner implements BucketAssigner<InputRecordPojo, String> {
@Override
public String getBucketId
(InputRecordPojo record,
Context context
) { return extractCustomerID(record);
}
@Override
public SimpleVersionedSerializer<String> getSerializer() {
return SimpleVersionedStringSerializer.INSTANCE;
}
}
FileSink<InputRecordPojo> sink =
FileSink.forRowFormat(...)
.withBucketAssigner(new FileBucketAssigner())
.build();
aW1wb3J0IG9yZy5hcGFjaGUuZmxpbmsuY29yZS5pby5TaW1wbGVWZXJzaW9uZWRTZXJpYWxpemVyOwppbXBvcnQgb3JnLmFwYWNoZS5mbGluay5zdHJlYW1pbmcuYXBpLmZ1bmN0aW9ucy5zaW5rLmZpbGVzeXN0ZW0uQnVja2V0QXNzaWduZXI7CmltcG9ydCBvcmcuYXBhY2hlLmZsaW5rLnN0cmVhbWluZy5hcGkuZnVuY3Rpb25zLnNpbmsuZmlsZXN5c3RlbS5idWNrZXRhc3NpZ25lcnMuU2ltcGxlVmVyc2lvbmVkU3RyaW5nU2VyaWFsaXplcjsKCnB1YmxpYyBjbGFzcyBGaWxlQnVja2V0QXNzaWduZXIgaW1wbGVtZW50cyBCdWNrZXRBc3NpZ25lcjxJbnB1dFJlY29yZFBvam8sIFN0cmluZz4gewoKICBAT3ZlcnJpZGUKICBwdWJsaWMgU3RyaW5nIGdldEJ1Y2tldElkKElucHV0UmVjb3JkUG9qbyByZWNvcmQsIENvbnRleHQgY29udGV4dCkgewoJcmV0dXJuIGV4dHJhY3RDdXN0b21lcklEKHJlY29yZCk7CiAgfQoKICBAT3ZlcnJpZGUKICBwdWJsaWMgU2ltcGxlVmVyc2lvbmVkU2VyaWFsaXplcjxTdHJpbmc+IGdldFNlcmlhbGl6ZXIoKSB7CiAgICByZXR1cm4gU2ltcGxlVmVyc2lvbmVkU3RyaW5nU2VyaWFsaXplci5JTlNUQU5DRTsKICB9Cn0KCkZpbGVTaW5rPElucHV0UmVjb3JkUG9qbz4gc2luayA9CiAgICAgICAgRmlsZVNpbmsuZm9yUm93Rm9ybWF0KC4uLikKICAgICAgICAgICAgLndpdGhCdWNrZXRBc3NpZ25lcihuZXcgRmlsZUJ1Y2tldEFzc2lnbmVyKCkpCiAgICAgICAgICAgIC5idWlsZCgpOw==