fork download
  1. import org.apache.flink.core.io.SimpleVersionedSerializer;
  2. import org.apache.flink.streaming.api.functions.sink.filesystem.BucketAssigner;
  3. import org.apache.flink.streaming.api.functions.sink.filesystem.bucketassigners.SimpleVersionedStringSerializer;
  4.  
  5. public class FileBucketAssigner implements BucketAssigner<InputRecordPojo, String> {
  6.  
  7. @Override
  8. public String getBucketId(InputRecordPojo record, Context context) {
  9. return extractCustomerID(record);
  10. }
  11.  
  12. @Override
  13. public SimpleVersionedSerializer<String> getSerializer() {
  14. return SimpleVersionedStringSerializer.INSTANCE;
  15. }
  16. }
  17.  
  18. FileSink<InputRecordPojo> sink =
  19. FileSink.forRowFormat(...)
  20. .withBucketAssigner(new FileBucketAssigner())
  21. .build();
Compilation error #stdin compilation error #stdout 0s 0KB
stdin
Standard input is empty
compilation info
Main.java:18: error: class, interface, or enum expected
FileSink<InputRecordPojo> sink =
^
1 error
stdout
Standard output is empty