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();