import { StorageType } from './../common/types'; import { MasterServer } from './MasterServer'; import { registerHandler } from '../common/handler'; import { Request } from '../client/Client'; import { SerializedFunction, deserialize, serialize, } from '../common/SerializeFunction'; export const CREATE_RDD = '@@master/createRDD'; export const MAP = '@@master/map'; export const REDUCE = '@@master/reduce'; export const REPARTITION = '@@master/repartition'; export const COALESCE = '@@master/coalesce'; export const CONCAT = '@@master/concat'; export const LOAD_FILE = '@@master/loadFile'; export const GET_NUM_PARTITIONS = '@@master/getNumPartitions'; export const SAVE_FILE = '@@master/saveFile'; export const SORT = '@@master/SORT'; export const CACHE = '@@master/cache'; export const LOAD_CACHE = '@@master/loadCache'; export const RELEASE_CACHE = '@@master/releaseCache'; registerHandler( REDUCE, async ( { subRequest, partitionFunc, finalFunc, }: { subRequest: Request; partitionFunc: SerializedFunction<(arg: T[]) => T1>; finalFunc: SerializedFunction<(arg: T1[]) => T1>; }, context: MasterServer, ) => { const results: T1[] = await context.runWork( subRequest, { type: 'reduce', }, [partitionFunc], ); return deserialize(finalFunc)(results); }, ); registerHandler( SAVE_FILE, async ( { subRequest, baseUrl, overwrite = true, serializer, extension = 'txt', }: { subRequest: Request; baseUrl: string; overwrite?: boolean; serializer: (data: any[]) => Buffer | Promise; extension?: string; }, context: MasterServer, ) => { const fileLoader = await context.getFileLoader(baseUrl, 'save'); await fileLoader.initSaveProgress(baseUrl, overwrite); const numPartitions = await context.getPartitionCount(subRequest); const radix = numPartitions <= 10000000 ? 10 : numPartitions <= 0x10000000 ? 16 : 36; const digits = numPartitions <= 100000 ? 5 : numPartitions <= 1000000 ? 6 : 7; const args = new Array(numPartitions) .fill(0) .map( (v, i) => `part-${('0000000' + i.toString(radix)).substr( -digits, )}.${extension}`, ); const saver = fileLoader.createDataSaver(baseUrl); const saveFunc = serialize( async (data: any[], filename: string) => { const buffer = await serializer(data); return saver(filename, buffer); }, { saver, serializer, }, ); await context.runWork(subRequest, { type: 'saveFile', saveFunc, args, }); await fileLoader.markSaveSuccess(baseUrl); }, ); registerHandler( GET_NUM_PARTITIONS, (subRequest: Request, context: MasterServer) => { return context.getPartitionCount(subRequest); }, ); registerHandler( CACHE, async ( { subRequest, storageType, }: { subRequest: Request; storageType: StorageType; }, context: MasterServer, ) => { const partitions = await context.runWork(subRequest, { type: 'partitions', storageType, }); return context.addCache(storageType, partitions); }, ); registerHandler(RELEASE_CACHE, (id: number, context: MasterServer) => { return context.releaseCache(id); });