| 30 | TASK_MEM = 128 |
| 31 | |
| 32 | class TestScheduler(mesos.interface.Scheduler): |
| 33 | def __init__(self, implicitAcknowledgements, executor, framework): |
| 34 | self.implicitAcknowledgements = implicitAcknowledgements |
| 35 | self.executor = executor |
| 36 | self.framework = framework |
| 37 | self.taskData = {} |
| 38 | self.tasksLaunched = 0 |
| 39 | self.tasksFinished = 0 |
| 40 | self.messagesSent = 0 |
| 41 | self.messagesReceived = 0 |
| 42 | |
| 43 | def registered(self, driver, frameworkId, masterInfo): |
| 44 | print "Registered with framework ID %s" % frameworkId.value |
| 45 | self.framework.id.CopyFrom(frameworkId) |
| 46 | driver.updateFramework(framework, []) |
| 47 | |
| 48 | def resourceOffers(self, driver, offers): |
| 49 | for offer in offers: |
| 50 | tasks = [] |
| 51 | offerCpus = 0 |
| 52 | offerMem = 0 |
| 53 | for resource in offer.resources: |
| 54 | if resource.name == "cpus": |
| 55 | offerCpus += resource.scalar.value |
| 56 | elif resource.name == "mem": |
| 57 | offerMem += resource.scalar.value |
| 58 | |
| 59 | print "Received offer %s with cpus: %s and mem: %s" \ |
| 60 | % (offer.id.value, offerCpus, offerMem) |
| 61 | |
| 62 | remainingCpus = offerCpus |
| 63 | remainingMem = offerMem |
| 64 | |
| 65 | while self.tasksLaunched < TOTAL_TASKS and \ |
| 66 | remainingCpus >= TASK_CPUS and \ |
| 67 | remainingMem >= TASK_MEM: |
| 68 | tid = self.tasksLaunched |
| 69 | self.tasksLaunched += 1 |
| 70 | |
| 71 | print "Launching task %d using offer %s" \ |
| 72 | % (tid, offer.id.value) |
| 73 | |
| 74 | task = mesos_pb2.TaskInfo() |
| 75 | task.task_id.value = str(tid) |
| 76 | task.slave_id.value = offer.slave_id.value |
| 77 | task.name = "task %d" % tid |
| 78 | task.executor.MergeFrom(self.executor) |
| 79 | |
| 80 | cpus = task.resources.add() |
| 81 | cpus.name = "cpus" |
| 82 | cpus.type = mesos_pb2.Value.SCALAR |
| 83 | cpus.scalar.value = TASK_CPUS |
| 84 | |
| 85 | mem = task.resources.add() |
| 86 | mem.name = "mem" |
| 87 | mem.type = mesos_pb2.Value.SCALAR |
| 88 | mem.scalar.value = TASK_MEM |
| 89 | |