blob: c959d5dec73df805c1559ba6e6a949dc4e7b3ea4 (
plain) (
blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
|
"""
Unit tests for PySpark; additional tests are implemented as doctests in
individual modules.
"""
import atexit
import os
import shutil
from tempfile import NamedTemporaryFile
import time
import unittest
from pyspark.context import SparkContext
class TestCheckpoint(unittest.TestCase):
def setUp(self):
self.sc = SparkContext('local[4]', 'TestPartitioning', batchSize=2)
def tearDown(self):
self.sc.stop()
def test_basic_checkpointing(self):
checkpointDir = NamedTemporaryFile(delete=False)
os.unlink(checkpointDir.name)
self.sc.setCheckpointDir(checkpointDir.name)
parCollection = self.sc.parallelize([1, 2, 3, 4])
flatMappedRDD = parCollection.flatMap(lambda x: range(1, x + 1))
self.assertFalse(flatMappedRDD.isCheckpointed())
self.assertIsNone(flatMappedRDD.getCheckpointFile())
flatMappedRDD.checkpoint()
result = flatMappedRDD.collect()
time.sleep(1) # 1 second
self.assertTrue(flatMappedRDD.isCheckpointed())
self.assertEqual(flatMappedRDD.collect(), result)
self.assertEqual(checkpointDir.name,
os.path.dirname(flatMappedRDD.getCheckpointFile()))
atexit.register(lambda: shutil.rmtree(checkpointDir.name))
if __name__ == "__main__":
unittest.main()
|